Compare commits

..

2 commits

Author SHA1 Message Date
26d5805747
f-33: vm: start vms with multiple nic
All checks were successful
Release Pipeline / checksums (push) Successful in 4s
Release Pipeline / release (push) Successful in 11s
Release Pipeline / set-release-target (push) Successful in 2s
Release Pipeline / upload-assets (agent.service, systemd/agent.service) (push) Successful in 4s
Release Pipeline / upload-assets (dnsmasq@.service, systemd/dnsmasq@.service) (push) Successful in 4s
Release Pipeline / upload-assets (metadata@.service, systemd/metadata@.service) (push) Successful in 4s
Release Pipeline / upload-assets (run-dnsmasq-in-netns.sh, scripts/run-dnsmasq-in-netns.sh) (push) Successful in 4s
Release Pipeline / build (agent, amd64, linux) (push) Successful in 0s
Release Pipeline / build (metadata, amd64, linux) (push) Successful in 0s
Release Pipeline / publish (push) Successful in 0s
Release Pipeline / build (push) Successful in 1m36s
Signed-off-by: GnomeZworc <nicolas.boufidjeline@g3e.fr>
2026-08-25 23:12:55 +02:00
0237acba81
f-33: net: add multi net in dispatch et vms #33
Signed-off-by: GnomeZworc <nicolas.boufidjeline@g3e.fr>
2026-08-25 22:48:13 +02:00
19 changed files with 1051 additions and 160 deletions

View file

@ -448,6 +448,12 @@ components:
interfaces: interfaces:
type: array type: array
minItems: 1 minItems: 1
description: >
Network interfaces, in order. The position determines the PCI slot
(0x03 + index) and therefore the interface name inside the guest.
**Exactly one** interface must be marked primary: it carries the
default route and the metadata server. All subnets must belong to
the same VPC.
items: items:
$ref: "#/components/schemas/VMInterface" $ref: "#/components/schemas/VMInterface"
storage: storage:

View file

@ -53,6 +53,11 @@ func main() {
return return
} }
if err := migration.MigrateVMNICs(db, log.With(slog.String("component", "migration"))); err != nil {
log.Error("vm nic migration failed", "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)

View file

@ -3,6 +3,7 @@ package agentapi
import ( import (
"encoding/json" "encoding/json"
"net/http" "net/http"
"sort"
"strconv" "strconv"
"strings" "strings"
@ -71,6 +72,52 @@ func (s *Server) stopVM(w http.ResponseWriter, _ *http.Request, name string) {
json.NewEncoder(w).Encode(vm) json.NewEncoder(w).Encode(vm)
} }
// interfacesFromDB reconstruit les interfaces depuis vm/<name>/nic/<index>/…,
// triées par index — celui-ci détermine le slot PCI, donc le nom de l'interface
// dans le guest.
func interfacesFromDB(prefix string, entries map[string]string) []VMInterface {
nicPrefix := prefix + "nic/"
byIndex := make(map[int]*VMInterface)
for key, value := range entries {
rest := strings.TrimPrefix(key, nicPrefix)
if rest == key {
continue
}
parts := strings.Split(rest, "/")
if len(parts) != 2 {
continue
}
idx, err := strconv.Atoi(parts[0])
if err != nil {
continue
}
if byIndex[idx] == nil {
byIndex[idx] = &VMInterface{}
}
switch parts[1] {
case "subnet":
byIndex[idx].Subnet = value
case "ip":
byIndex[idx].IP = value
case "primary":
byIndex[idx].Primary = value == "true"
}
}
indexes := make([]int, 0, len(byIndex))
for idx := range byIndex {
indexes = append(indexes, idx)
}
sort.Ints(indexes)
ifaces := make([]VMInterface, 0, len(indexes))
for _, idx := range indexes {
ifaces = append(ifaces, *byIndex[idx])
}
return ifaces
}
func vmFromDB(name string, entries map[string]string) (VM, error) { func vmFromDB(name string, entries map[string]string) (VM, error) {
prefix := "vm/" + name + "/" prefix := "vm/" + name + "/"
vm := VM{Name: name} vm := VM{Name: name}
@ -81,11 +128,7 @@ func vmFromDB(name string, entries map[string]string) (VM, error) {
vm.CPUs, _ = strconv.Atoi(entries[prefix+"cpus"]) vm.CPUs, _ = strconv.Atoi(entries[prefix+"cpus"])
vm.UEFI = entries[prefix+"uefi"] == "true" vm.UEFI = entries[prefix+"uefi"] == "true"
subnet := entries[prefix+"subnet"] vm.Interfaces = interfacesFromDB(prefix, entries)
ip := entries[prefix+"ip"]
if subnet != "" || ip != "" {
vm.Interfaces = []VMInterface{{Subnet: subnet, IP: ip, Primary: true}}
}
diskPrefix := prefix + "disk/" diskPrefix := prefix + "disk/"
for key, path := range entries { for key, path := range entries {

View file

@ -262,3 +262,89 @@ func TestStartVM_EmptyBase64MeansNoDocument(t *testing.T) {
t.Errorf("un base64 vide est indiscernable d'un champ absent : %v", entries) t.Errorf("un base64 vide est indiscernable d'un champ absent : %v", entries)
} }
} }
// --- interfaces multiples ---
func TestVmFromDB_MultipleInterfacesSortedByIndex(t *testing.T) {
vm, err := vmFromDB("vm-multi", map[string]string{
"vm/vm-multi/state": "running",
"vm/vm-multi/nic/1/subnet": "sn-2",
"vm/vm-multi/nic/1/ip": "10.2.0.5",
"vm/vm-multi/nic/0/subnet": "sn-1",
"vm/vm-multi/nic/0/ip": "10.1.0.5",
"vm/vm-multi/nic/0/primary": "true",
"vm/vm-multi/disk/vda": "/data/root.qcow2",
})
if err != nil {
t.Fatalf("vmFromDB : %v", err)
}
if len(vm.Interfaces) != 2 {
t.Fatalf("2 interfaces attendues, obtenu %d : %+v", len(vm.Interfaces), vm.Interfaces)
}
if vm.Interfaces[0].Subnet != "sn-1" || !vm.Interfaces[0].Primary {
t.Errorf("la première doit être l'index 0, primaire : %+v", vm.Interfaces[0])
}
if vm.Interfaces[1].Subnet != "sn-2" || vm.Interfaces[1].Primary {
t.Errorf("la seconde doit être l'index 1, non primaire : %+v", vm.Interfaces[1])
}
}
func TestStartVM_StoresAllInterfaces(t *testing.T) {
s, db := newTestServer(t)
for _, sn := range []string{"sn-1", "sn-2"} {
kv.AddInDB(db, "subnet/"+sn+"/state", "running")
kv.AddInDB(db, "subnet/"+sn+"/vpc", "vpc-1")
}
body, _ := json.Marshal(VMCreateRequest{
Name: "vm-multi",
Interfaces: []VMInterface{
{Subnet: "sn-1", IP: "10.1.0.5", Primary: true},
{Subnet: "sn-2", IP: "10.2.0.5"},
},
Storage: []VMStorage{{Path: "/data/root.qcow2", Dev: "vda"}},
})
w := httptest.NewRecorder()
s.VmsHandler(w, httptest.NewRequest(http.MethodPost, "/vms", bytes.NewReader(body)))
if w.Code != http.StatusAccepted {
t.Fatalf("attendu 202, obtenu %d : %s", w.Code, w.Body.String())
}
if got, _ := kv.GetFromDB(db, "vm/vm-multi/nic/1/subnet"); got != "sn-2" {
t.Errorf("seconde interface non stockée : %q", got)
}
if got, _ := kv.GetFromDB(db, "vm/vm-multi/nic/0/primary"); got != "true" {
t.Errorf("primaire non marquée : %q", got)
}
if _, err := kv.GetFromDB(db, "vm/vm-multi/nic/1/primary"); err == nil {
t.Error("une interface non primaire ne doit pas porter la clé primary")
}
}
func TestStartVM_RejectsZeroOrTwoPrimaries(t *testing.T) {
cases := map[string][]VMInterface{
"aucune primaire": {{Subnet: "sn-1", IP: "10.1.0.5"}},
"deux primaires": {
{Subnet: "sn-1", IP: "10.1.0.5", Primary: true},
{Subnet: "sn-2", IP: "10.2.0.5", Primary: true},
},
}
for label, ifaces := range cases {
s, db := newTestServer(t)
for _, sn := range []string{"sn-1", "sn-2"} {
kv.AddInDB(db, "subnet/"+sn+"/state", "running")
kv.AddInDB(db, "subnet/"+sn+"/vpc", "vpc-1")
}
body, _ := json.Marshal(VMCreateRequest{
Name: "vm-bad",
Interfaces: ifaces,
Storage: []VMStorage{{Path: "/data/root.qcow2", Dev: "vda"}},
})
w := httptest.NewRecorder()
s.VmsHandler(w, httptest.NewRequest(http.MethodPost, "/vms", bytes.NewReader(body)))
if w.Code != http.StatusBadRequest {
t.Errorf("%s : attendu 400, obtenu %d — %s", label, w.Code, w.Body.String())
}
}
}

View file

@ -67,16 +67,17 @@ func (s *Server) startVM(w http.ResponseWriter, r *http.Request) {
return return
} }
var primary *VMInterface nics := make([]dispatcher.VMNIC, len(req.Interfaces))
for i := range req.Interfaces { primaries := 0
if req.Interfaces[i].Primary { for i, iface := range req.Interfaces {
primary = &req.Interfaces[i] nics[i] = dispatcher.VMNIC{Subnet: iface.Subnet, IP: iface.IP, Primary: iface.Primary}
break if iface.Primary {
primaries++
} }
} }
if primary == nil { if primaries != 1 {
w.WriteHeader(http.StatusBadRequest) w.WriteHeader(http.StatusBadRequest)
json.NewEncoder(w).Encode(ErrorResponse{Error: "one interface must be primary"}) json.NewEncoder(w).Encode(ErrorResponse{Error: "exactly one interface must be primary"})
return return
} }
@ -94,8 +95,7 @@ func (s *Server) startVM(w http.ResponseWriter, r *http.Request) {
cmd := dispatcher.StartVMCommand{ cmd := dispatcher.StartVMCommand{
Name: req.Name, Name: req.Name,
Subnet: primary.Subnet, NICs: nics,
IP: primary.IP,
Disks: disks, Disks: disks,
Memory: req.Memory, Memory: req.Memory,
CPUs: req.CPUs, CPUs: req.CPUs,

View file

@ -10,6 +10,7 @@ import (
type Reservation struct { type Reservation struct {
MAC string MAC string
IP string IP string
Tag string // pose set:<Tag> sur l'entrée, pour cibler les options par interface
} }
func HostsDir(confDir, name string) string { func HostsDir(confDir, name string) string {
@ -39,7 +40,11 @@ func WriteReservations(confDir, name, vmName string, res []Reservation) error {
if r.MAC == "" || r.IP == "" { if r.MAC == "" || r.IP == "" {
return fmt.Errorf("incomplete reservation for vm %q: mac=%q ip=%q", vmName, r.MAC, r.IP) return fmt.Errorf("incomplete reservation for vm %q: mac=%q ip=%q", vmName, r.MAC, r.IP)
} }
if r.Tag == "" {
fmt.Fprintf(&sb, "%s,%s\n", r.MAC, r.IP) fmt.Fprintf(&sb, "%s,%s\n", r.MAC, r.IP)
continue
}
fmt.Fprintf(&sb, "%s,%s,set:%s\n", r.MAC, r.IP, r.Tag)
} }
path := filepath.Join(dir, vmName) path := filepath.Join(dir, vmName)
@ -49,6 +54,51 @@ func WriteReservations(confDir, name, vmName string, res []Reservation) error {
return nil return nil
} }
// WriteVMOptions écrit les options DHCP propres à des interfaces de cette VM
// sur ce subnet : elles suppriment la route par défaut — option 3 nue — et
// réémettent les autres routes.
//
// Un override de l'option 121 remplace la précédente en entier, il ne s'y
// ajoute pas (vérifié sur dnsmasq 2.90) : omettre la route vers le serveur de
// métadonnées la ferait disparaître, et la VM ne se provisionnerait pas.
func WriteVMOptions(confDir, name, vmName string, tags []string, c Config) error {
if len(tags) == 0 {
return RemoveVMOptions(confDir, name, vmName)
}
if c.InterfaceIP == nil {
return fmt.Errorf("interface ip is required to build options for vm %q", vmName)
}
if c.DefaultGateway != nil {
return fmt.Errorf("vm options for %q must not carry a default route", vmName)
}
dir := OptsDir(confDir, name)
if err := os.MkdirAll(dir, 0755); err != nil {
return fmt.Errorf("create %s: %w", dir, err)
}
routes := strings.Join(classlessRoutes(c), ",")
var sb strings.Builder
for _, tag := range tags {
fmt.Fprintf(&sb, "tag:%s,3\n", tag)
fmt.Fprintf(&sb, "tag:%s,121,%s\n", tag, routes)
}
path := filepath.Join(dir, vmName)
if err := os.WriteFile(path, []byte(sb.String()), 0644); err != nil {
return fmt.Errorf("write %s: %w", path, err)
}
return nil
}
func RemoveVMOptions(confDir, name, vmName string) error {
path := filepath.Join(OptsDir(confDir, name), vmName)
if err := os.Remove(path); err != nil && !os.IsNotExist(err) {
return fmt.Errorf("remove %s: %w", path, err)
}
return nil
}
func RemoveReservations(confDir, name, vmName string) error { func RemoveReservations(confDir, name, vmName string) error {
for _, dir := range []string{HostsDir(confDir, name), OptsDir(confDir, name)} { for _, dir := range []string{HostsDir(confDir, name), OptsDir(confDir, name)} {
path := filepath.Join(dir, vmName) path := filepath.Join(dir, vmName)

View file

@ -1,8 +1,10 @@
package dhcp package dhcp
import ( import (
"net"
"os" "os"
"path/filepath" "path/filepath"
"strings"
"testing" "testing"
) )
@ -124,3 +126,103 @@ func TestUnitName(t *testing.T) {
t.Errorf("unit attendue dnsmasq@%s.service, obtenu %s", subName, got) t.Errorf("unit attendue dnsmasq@%s.service, obtenu %s", subName, got)
} }
} }
// --- options par interface ---
func TestWriteReservations_WithTags(t *testing.T) {
dir := t.TempDir()
res := []Reservation{
{MAC: "00:22:33:00:01:02", IP: "10.1.1.2", Tag: "i-web-0"},
{MAC: "00:22:33:00:02:07", IP: "10.1.2.7", Tag: "i-web-1"},
}
if err := WriteReservations(dir, subName, "i-web", res); err != nil {
t.Fatalf("WriteReservations : %v", err)
}
content, _ := os.ReadFile(filepath.Join(HostsDir(dir, subName), "i-web"))
want := "00:22:33:00:01:02,10.1.1.2,set:i-web-0\n00:22:33:00:02:07,10.1.2.7,set:i-web-1\n"
if string(content) != want {
t.Errorf("attendu %q, obtenu %q", want, content)
}
}
func vmOptions(t *testing.T, dir string, tags []string, vpcRoute *net.IPNet) string {
t.Helper()
if err := WriteVMOptions(dir, subName, "i-web", tags, Config{
InterfaceIP: net.ParseIP("10.1.1.1").To4(),
VPCRoute: vpcRoute,
}); err != nil {
t.Fatalf("WriteVMOptions : %v", err)
}
b, err := os.ReadFile(filepath.Join(OptsDir(dir, subName), "i-web"))
if err != nil {
return ""
}
return string(b)
}
func TestWriteVMOptions_SuppressesDefaultRoute(t *testing.T) {
_, vpcNet, _ := net.ParseCIDR("192.168.0.0/16")
got := vmOptions(t, t.TempDir(), []string{"i-web-1"}, vpcNet)
if !strings.Contains(got, "tag:i-web-1,3\n") {
t.Errorf("l'option 3 nue doit supprimer la route par défaut :\n%s", got)
}
if strings.Contains(got, "0.0.0.0/0") {
t.Errorf("la route par défaut ne doit pas figurer dans l'override :\n%s", got)
}
}
func TestWriteVMOptions_ReEmitsMetadataRoute(t *testing.T) {
_, vpcNet, _ := net.ParseCIDR("192.168.0.0/16")
got := vmOptions(t, t.TempDir(), []string{"i-web-1"}, vpcNet)
if !strings.Contains(got, "169.254.169.254/32,10.1.1.1") {
t.Errorf("un override de l'option 121 remplace la précédente en entier : sans la route metadata, la VM ne se provisionne pas\n%s", got)
}
if !strings.Contains(got, "192.168.0.0/16,10.1.1.1") {
t.Errorf("la route VPC doit être réémise elle aussi :\n%s", got)
}
}
func TestWriteVMOptions_BridgeHasNoVPCRoute(t *testing.T) {
got := vmOptions(t, t.TempDir(), []string{"i-web-1"}, nil)
if !strings.Contains(got, "169.254.169.254/32") {
t.Errorf("route metadata absente :\n%s", got)
}
if strings.Contains(got, "192.168") {
t.Errorf("aucune route VPC attendue en mode bridge :\n%s", got)
}
}
func TestWriteVMOptions_OneBlockPerTag(t *testing.T) {
got := vmOptions(t, t.TempDir(), []string{"i-web-1", "i-web-2"}, nil)
for _, tag := range []string{"tag:i-web-1,3", "tag:i-web-2,3"} {
if !strings.Contains(got, tag) {
t.Errorf("%q absent :\n%s", tag, got)
}
}
}
func TestWriteVMOptions_NoTagRemovesFile(t *testing.T) {
dir := t.TempDir()
_ = vmOptions(t, dir, []string{"i-web-1"}, nil)
if err := WriteVMOptions(dir, subName, "i-web", nil, Config{}); err != nil {
t.Fatalf("WriteVMOptions : %v", err)
}
if _, err := os.Stat(filepath.Join(OptsDir(dir, subName), "i-web")); !os.IsNotExist(err) {
t.Error("sans interface non primaire, aucun fichier d'options ne doit subsister")
}
}
func TestWriteVMOptions_RefusesDefaultGateway(t *testing.T) {
err := WriteVMOptions(t.TempDir(), subName, "i-web", []string{"i-web-1"}, Config{
InterfaceIP: net.ParseIP("10.1.1.1").To4(),
DefaultGateway: net.ParseIP("10.1.1.254").To4(),
})
if err == nil {
t.Error("ces options servent à retirer la route par défaut : en porter une est une incohérence")
}
}

View file

@ -19,10 +19,15 @@ type VMDisk struct {
Dev string Dev string
} }
type StartVMCommand struct { type VMNIC struct {
Name string
Subnet string Subnet string
IP string IP string
Primary bool
}
type StartVMCommand struct {
Name string
NICs []VMNIC
Disks []VMDisk Disks []VMDisk
Memory int Memory int
CPUs int CPUs int
@ -38,20 +43,28 @@ func (c StartVMCommand) Prepare(db *badger.DB, _ *configuration.Config) error {
if _, err := kv.GetFromDB(db, "vm/"+c.Name+"/state"); err == nil { if _, err := kv.GetFromDB(db, "vm/"+c.Name+"/state"); err == nil {
return fmt.Errorf("vm %q already exists", c.Name) return fmt.Errorf("vm %q already exists", c.Name)
} }
subnetState, err := state.Get(db, "subnet/"+c.Subnet) if err := c.validateNICs(db); err != nil {
if err != nil { return err
return fmt.Errorf("subnet %q not found", c.Subnet)
}
if subnetState != state.Creating && subnetState != state.Running {
return fmt.Errorf("subnet %q is %s", c.Subnet, subnetState)
} }
port, err := allocateMetadataPort(db) port, err := allocateMetadataPort(db)
if err != nil { if err != nil {
return fmt.Errorf("allocate metadata port: %w", err) return fmt.Errorf("allocate metadata port: %w", err)
} }
state.Set(db, c.Key(), state.Creating) state.Set(db, c.Key(), state.Creating)
kv.AddInDB(db, "vm/"+c.Name+"/subnet", c.Subnet) for i, n := range c.NICs {
kv.AddInDB(db, "vm/"+c.Name+"/ip", c.IP) prefix := fmt.Sprintf("vm/%s/nic/%d/", c.Name, i)
if err := kv.AddInDB(db, prefix+"subnet", n.Subnet); err != nil {
return fmt.Errorf("store nic %d subnet: %w", i, err)
}
if err := kv.AddInDB(db, prefix+"ip", n.IP); err != nil {
return fmt.Errorf("store nic %d ip: %w", i, err)
}
if n.Primary {
if err := kv.AddInDB(db, prefix+"primary", "true"); err != nil {
return fmt.Errorf("store nic %d primary: %w", i, err)
}
}
}
kv.AddInDB(db, "vm/"+c.Name+"/metadata_port", strconv.Itoa(port)) kv.AddInDB(db, "vm/"+c.Name+"/metadata_port", strconv.Itoa(port))
for _, d := range c.Disks { for _, d := range c.Disks {
kv.AddInDB(db, "vm/"+c.Name+"/disk/"+d.Dev, d.Path) kv.AddInDB(db, "vm/"+c.Name+"/disk/"+d.Dev, d.Path)
@ -75,6 +88,31 @@ func (c StartVMCommand) Prepare(db *badger.DB, _ *configuration.Config) error {
return nil return nil
} }
// validateNICs vérifie qu'il y a exactement une interface primaire et que
// chaque subnet référencé est utilisable.
func (c StartVMCommand) validateNICs(db *badger.DB) error {
if len(c.NICs) == 0 {
return fmt.Errorf("vm %q has no interface", c.Name)
}
primaries := 0
for _, n := range c.NICs {
if n.Primary {
primaries++
}
subnetState, err := state.Get(db, "subnet/"+n.Subnet)
if err != nil {
return fmt.Errorf("subnet %q not found", n.Subnet)
}
if subnetState != state.Creating && subnetState != state.Running {
return fmt.Errorf("subnet %q is %s", n.Subnet, subnetState)
}
}
if primaries != 1 {
return fmt.Errorf("vm %q has %d primary interfaces, expected exactly one", c.Name, primaries)
}
return nil
}
func allocateMetadataPort(db *badger.DB) (int, error) { func allocateMetadataPort(db *badger.DB) (int, error) {
entries, err := kv.ListByPrefix(db, "vm/") entries, err := kv.ListByPrefix(db, "vm/")
if err != nil { if err != nil {
@ -99,23 +137,25 @@ func allocateMetadataPort(db *badger.DB) (int, error) {
func (c StartVMCommand) Execute(db *badger.DB, cfg *configuration.Config) error { func (c StartVMCommand) Execute(db *badger.DB, cfg *configuration.Config) error {
timeout := time.After(time.Duration(cfg.Dispatcher.TimeoutSeconds) * time.Second) timeout := time.After(time.Duration(cfg.Dispatcher.TimeoutSeconds) * time.Second)
for _, n := range c.NICs {
for { for {
subnetState, err := state.Get(db, "subnet/"+c.Subnet) subnetState, err := state.Get(db, "subnet/"+n.Subnet)
if err != nil { if err != nil {
return fmt.Errorf("subnet %q not found while waiting", c.Subnet) return fmt.Errorf("subnet %q not found while waiting", n.Subnet)
} }
if subnetState == state.Running { if subnetState == state.Running {
break break
} }
if subnetState != state.Creating { if subnetState != state.Creating {
return fmt.Errorf("subnet %q is %s, cannot start vm %q", c.Subnet, subnetState, c.Name) return fmt.Errorf("subnet %q is %s, cannot start vm %q", n.Subnet, subnetState, c.Name)
} }
select { select {
case <-timeout: case <-timeout:
return fmt.Errorf("timed out waiting for subnet %q to be running", c.Subnet) return fmt.Errorf("timed out waiting for subnet %q to be running", n.Subnet)
case <-time.After(time.Duration(cfg.Dispatcher.PollSeconds) * time.Second): case <-time.After(time.Duration(cfg.Dispatcher.PollSeconds) * time.Second):
} }
} }
}
return vm.StartVM(db, c.Name, cfg) return vm.StartVM(db, c.Name, cfg)
} }

View file

@ -16,8 +16,7 @@ func TestStartVMCommand_Prepare_SingleDisk(t *testing.T) {
cmd := StartVMCommand{ cmd := StartVMCommand{
Name: "vm-1", Name: "vm-1",
Subnet: "sn-1", NICs: []VMNIC{{Subnet: "sn-1", IP: "10.0.0.5", Primary: true}},
IP: "10.0.0.5",
Disks: []VMDisk{{Path: "/data/root.qcow2", Dev: "sda"}}, Disks: []VMDisk{{Path: "/data/root.qcow2", Dev: "sda"}},
} }
if err := cmd.Prepare(db, nil); err != nil { if err := cmd.Prepare(db, nil); err != nil {
@ -40,8 +39,7 @@ func TestStartVMCommand_Prepare_MultiDisk(t *testing.T) {
cmd := StartVMCommand{ cmd := StartVMCommand{
Name: "vm-2", Name: "vm-2",
Subnet: "sn-1", NICs: []VMNIC{{Subnet: "sn-1", IP: "10.0.0.6", Primary: true}},
IP: "10.0.0.6",
Disks: []VMDisk{ Disks: []VMDisk{
{Path: "/data/root.qcow2", Dev: "sda"}, {Path: "/data/root.qcow2", Dev: "sda"},
{Path: "/data/data.qcow2", Dev: "sdb"}, {Path: "/data/data.qcow2", Dev: "sdb"},
@ -73,8 +71,7 @@ func TestStartVMCommand_Prepare_SlotGap(t *testing.T) {
// sdb absent au boot — slot réservé pour hotplug // sdb absent au boot — slot réservé pour hotplug
cmd := StartVMCommand{ cmd := StartVMCommand{
Name: "vm-3", Name: "vm-3",
Subnet: "sn-1", NICs: []VMNIC{{Subnet: "sn-1", IP: "10.0.0.7", Primary: true}},
IP: "10.0.0.7",
Disks: []VMDisk{ Disks: []VMDisk{
{Path: "/data/root.qcow2", Dev: "sda"}, {Path: "/data/root.qcow2", Dev: "sda"},
{Path: "/data/extra.qcow2", Dev: "sdc"}, {Path: "/data/extra.qcow2", Dev: "sdc"},
@ -102,8 +99,7 @@ func TestStartVMCommand_Prepare_NoVolumePath(t *testing.T) {
cmd := StartVMCommand{ cmd := StartVMCommand{
Name: "vm-4", Name: "vm-4",
Subnet: "sn-1", NICs: []VMNIC{{Subnet: "sn-1", IP: "10.0.0.8", Primary: true}},
IP: "10.0.0.8",
Disks: []VMDisk{{Path: "/data/root.qcow2", Dev: "sda"}}, Disks: []VMDisk{{Path: "/data/root.qcow2", Dev: "sda"}},
} }
if err := cmd.Prepare(db, nil); err != nil { if err := cmd.Prepare(db, nil); err != nil {
@ -166,8 +162,7 @@ func TestStartVMCommand_Prepare_Duplicate(t *testing.T) {
cmd := StartVMCommand{ cmd := StartVMCommand{
Name: "vm-exist", Name: "vm-exist",
Subnet: "sn-1", NICs: []VMNIC{{Subnet: "sn-1", IP: "10.0.0.9", Primary: true}},
IP: "10.0.0.9",
Disks: []VMDisk{{Path: "/data/root.qcow2", Dev: "sda"}}, Disks: []VMDisk{{Path: "/data/root.qcow2", Dev: "sda"}},
} }
if err := cmd.Prepare(db, nil); err == nil { if err := cmd.Prepare(db, nil); err == nil {
@ -185,8 +180,7 @@ func prepareWithDocuments(t *testing.T, docs map[string]string) *badger.DB {
cmd := StartVMCommand{ cmd := StartVMCommand{
Name: "vm-doc", Name: "vm-doc",
Subnet: "sn-1", NICs: []VMNIC{{Subnet: "sn-1", IP: "10.0.0.5", Primary: true}},
IP: "10.0.0.5",
Disks: []VMDisk{{Path: "/data/root.qcow2", Dev: "vda"}}, Disks: []VMDisk{{Path: "/data/root.qcow2", Dev: "vda"}},
Documents: docs, Documents: docs,
} }

76
internal/migration/nic.go Normal file
View file

@ -0,0 +1,76 @@
package migration
import (
"fmt"
"log/slog"
"strings"
"git.g3e.fr/syonad/two/pkg/db/kv"
"github.com/dgraph-io/badger/v4"
)
// legacyNICKeys sont les clés d'interface d'avant le multi-subnet, portées
// directement par la VM.
var legacyNICKeys = []string{"subnet", "ip", "tap_id"}
// MigrateVMNICs déplace les clés d'interface de vm/<name>/<clé> vers
// vm/<name>/nic/0/<clé> et marque cette interface comme primaire.
//
// Idempotente : une VM possédant déjà des clés nic/ est laissée telle quelle.
// Sans cette migration, toute VM créée avant le passage au multi-subnet
// deviendrait illisible par loadVM.
func MigrateVMNICs(db *badger.DB, log *slog.Logger) error {
entries, err := kv.ListByPrefix(db, "vm/")
if err != nil {
return fmt.Errorf("list vm/: %w", err)
}
for _, name := range vmsToMigrate(entries) {
for _, key := range legacyNICKeys {
value, ok := entries["vm/"+name+"/"+key]
if !ok {
continue
}
if err := kv.AddInDB(db, "vm/"+name+"/nic/0/"+key, value); err != nil {
return fmt.Errorf("migrate %s of vm %s: %w", key, name, err)
}
if err := kv.DeleteInDB(db, "vm/"+name+"/"+key); err != nil {
return fmt.Errorf("delete legacy %s of vm %s: %w", key, name, err)
}
}
if err := kv.AddInDB(db, "vm/"+name+"/nic/0/primary", "true"); err != nil {
return fmt.Errorf("mark nic 0 primary for vm %s: %w", name, err)
}
log.Info("vm nics migrated", "resource", "vm/"+name, "reason", "legacy single interface")
}
return nil
}
// vmsToMigrate retourne les VM portant l'ancien schéma et aucune clé nic/.
func vmsToMigrate(entries map[string]string) []string {
legacy := make(map[string]bool)
migrated := make(map[string]bool)
for key := range entries {
parts := strings.Split(key, "/")
if len(parts) < 3 || parts[0] != "vm" {
continue
}
name := parts[1]
switch {
case len(parts) == 3 && (parts[2] == "subnet" || parts[2] == "ip" || parts[2] == "tap_id"):
legacy[name] = true
case parts[2] == "nic":
migrated[name] = true
}
}
var names []string
for name := range legacy {
if !migrated[name] {
names = append(names, name)
}
}
return names
}

View file

@ -0,0 +1,136 @@
package migration
import (
"io"
"log/slog"
"testing"
"git.g3e.fr/syonad/two/pkg/db/kv"
"github.com/dgraph-io/badger/v4"
)
func newNICDB(t *testing.T) *badger.DB {
t.Helper()
db := kv.InitDB(kv.Config{Path: t.TempDir()}, false)
t.Cleanup(func() { db.Close() })
return db
}
func quietLog() *slog.Logger {
return slog.New(slog.NewTextHandler(io.Discard, nil))
}
func seedLegacyVM(t *testing.T, db *badger.DB, name string) {
t.Helper()
kv.AddInDB(db, "vm/"+name+"/state", "running")
kv.AddInDB(db, "vm/"+name+"/subnet", "sn-000001")
kv.AddInDB(db, "vm/"+name+"/ip", "10.1.1.2")
kv.AddInDB(db, "vm/"+name+"/tap_id", "12345678")
}
func TestMigrateVMNICs_MovesLegacyKeys(t *testing.T) {
db := newNICDB(t)
seedLegacyVM(t, db, "i-test1")
if err := MigrateVMNICs(db, quietLog()); err != nil {
t.Fatalf("MigrateVMNICs : %v", err)
}
want := map[string]string{
"vm/i-test1/nic/0/subnet": "sn-000001",
"vm/i-test1/nic/0/ip": "10.1.1.2",
"vm/i-test1/nic/0/tap_id": "12345678",
"vm/i-test1/nic/0/primary": "true",
}
for key, expected := range want {
got, err := kv.GetFromDB(db, key)
if err != nil {
t.Errorf("clé %s absente : %v", key, err)
continue
}
if got != expected {
t.Errorf("%s = %q, attendu %q", key, got, expected)
}
}
}
func TestMigrateVMNICs_RemovesLegacyKeys(t *testing.T) {
db := newNICDB(t)
seedLegacyVM(t, db, "i-test1")
if err := MigrateVMNICs(db, quietLog()); err != nil {
t.Fatalf("MigrateVMNICs : %v", err)
}
for _, key := range []string{"vm/i-test1/subnet", "vm/i-test1/ip", "vm/i-test1/tap_id"} {
if _, err := kv.GetFromDB(db, key); err == nil {
t.Errorf("clé héritée %s toujours présente", key)
}
}
}
func TestMigrateVMNICs_PreservesOtherKeys(t *testing.T) {
db := newNICDB(t)
seedLegacyVM(t, db, "i-test1")
kv.AddInDB(db, "vm/i-test1/memory", "2048")
kv.AddInDB(db, "vm/i-test1/disk/vda", "/data/root.qcow2")
if err := MigrateVMNICs(db, quietLog()); err != nil {
t.Fatalf("MigrateVMNICs : %v", err)
}
if got, _ := kv.GetFromDB(db, "vm/i-test1/memory"); got != "2048" {
t.Errorf("memory altérée : %q", got)
}
if got, _ := kv.GetFromDB(db, "vm/i-test1/disk/vda"); got != "/data/root.qcow2" {
t.Errorf("disque altéré : %q", got)
}
}
func TestMigrateVMNICs_Idempotent(t *testing.T) {
db := newNICDB(t)
seedLegacyVM(t, db, "i-test1")
if err := MigrateVMNICs(db, quietLog()); err != nil {
t.Fatalf("premier passage : %v", err)
}
before, _ := kv.ListByPrefix(db, "vm/i-test1/")
if err := MigrateVMNICs(db, quietLog()); err != nil {
t.Fatalf("second passage : %v", err)
}
after, _ := kv.ListByPrefix(db, "vm/i-test1/")
if len(before) != len(after) {
t.Errorf("le second passage a modifié la base : %d clés puis %d", len(before), len(after))
}
for key, value := range before {
if after[key] != value {
t.Errorf("%s : %q devenu %q", key, value, after[key])
}
}
}
func TestMigrateVMNICs_LeavesMigratedVMsAlone(t *testing.T) {
db := newNICDB(t)
kv.AddInDB(db, "vm/i-multi/state", "running")
kv.AddInDB(db, "vm/i-multi/nic/0/subnet", "sn-000001")
kv.AddInDB(db, "vm/i-multi/nic/0/primary", "true")
kv.AddInDB(db, "vm/i-multi/nic/1/subnet", "sn-000002")
if err := MigrateVMNICs(db, quietLog()); err != nil {
t.Fatalf("MigrateVMNICs : %v", err)
}
if got, _ := kv.GetFromDB(db, "vm/i-multi/nic/1/subnet"); got != "sn-000002" {
t.Errorf("la seconde interface a été perdue : %q", got)
}
if _, err := kv.GetFromDB(db, "vm/i-multi/nic/0/primary"); err != nil {
t.Error("la primaire existante a été perdue")
}
}
func TestMigrateVMNICs_EmptyDB(t *testing.T) {
if err := MigrateVMNICs(newNICDB(t), quietLog()); err != nil {
t.Errorf("une base vide ne doit pas échouer : %v", err)
}
}

View file

@ -9,10 +9,16 @@ type DiskConfig struct {
Dev string Dev string
} }
type Config struct { // NICConfig décrit une interface réseau. Sa position dans Config.NICs
Name string // détermine le slot PCI, donc le nom de l'interface dans le guest.
type NICConfig struct {
TapID int TapID int
Mac string Mac string
}
type Config struct {
Name string
NICs []NICConfig
Disks []DiskConfig Disks []DiskConfig
Memory int Memory int
CPUs int CPUs int

View file

@ -11,6 +11,12 @@ import (
"strings" "strings"
) )
const (
firstNICSlot = 0x03
lastNICSlot = 0x1d
maxNICs = lastNICSlot - firstNICSlot + 1
)
func Start(cfg Config) error { func Start(cfg Config) error {
memory := cfg.Memory memory := cfg.Memory
if memory == 0 { if memory == 0 {
@ -97,11 +103,24 @@ func Start(cfg Config) error {
} }
} }
// Slots 0x03 à 0x1d réservés au réseau par la carte PCI (#36). Le slot est
// dérivé de l'index et non laissé à QEMU : c'est lui qui fixe le nom de
// l'interface dans le guest, et un slot flottant la renomme d'un démarrage
// à l'autre.
if len(cfg.NICs) == 0 {
return fmt.Errorf("vm %s has no network interface", cfg.Name)
}
if len(cfg.NICs) > maxNICs {
return fmt.Errorf("vm %s has %d interfaces, the pci map holds %d", cfg.Name, len(cfg.NICs), maxNICs)
}
for i, n := range cfg.NICs {
args = append(args, args = append(args,
"-netdev", fmt.Sprintf("tap,id=net0,ifname=tap%d,script=no,downscript=no", cfg.TapID), "-netdev", fmt.Sprintf("tap,id=net%d,ifname=tap%d,script=no,downscript=no", i, n.TapID),
"-device", fmt.Sprintf("virtio-net-pci,netdev=net0,mac=%s,bus=pci.0,addr=0x03", cfg.Mac), "-device", fmt.Sprintf("virtio-net-pci,netdev=net%d,mac=%s,bus=pci.0,addr=0x%02x", i, n.Mac, firstNICSlot+i),
"-daemonize",
) )
}
args = append(args, "-daemonize")
scopeArgs := append([]string{ scopeArgs := append([]string{
"--scope", "--scope",

View file

@ -3,6 +3,7 @@ package vm
import ( import (
"fmt" "fmt"
"io" "io"
"net"
"os" "os"
"path/filepath" "path/filepath"
@ -31,26 +32,36 @@ func StartVM(db *badger.DB, name string, cfg *configuration.Config) error {
if err != nil { if err != nil {
return err return err
} }
nic := d.primary()
if err := netif.CreateTap(d.tapID, d.bridge, d.vpcName); err != nil { for _, n := range d.nics {
return fmt.Errorf("create tap: %w", err) if err := netif.CreateTap(n.tapID, n.bridge, n.vpcName); err != nil {
return fmt.Errorf("create tap of interface %d: %w", n.index, err)
}
} }
if err := netns.Call(d.vpcName, func() error { // La redirection est posée pour chaque IP de la VM vers le serveur de
return iptables.AddMetadataRedirect(d.ip, d.interfaceIP, d.metadataPort) // métadonnées de l'interface primaire. Toutes les interfaces étant dans le
// même VPC, donc le même netns, il est joignable depuis n'importe laquelle.
if err := netns.Call(nic.vpcName, func() error {
for _, n := range d.nics {
if err := iptables.AddMetadataRedirect(n.ip, nic.interfaceIP, d.metadataPort); err != nil {
return fmt.Errorf("interface %d: %w", n.index, err)
}
}
return nil
}); err != nil { }); err != nil {
return fmt.Errorf("add metadata redirect: %w", err) return fmt.Errorf("add metadata redirect: %w", err)
} }
if err := dhcp.WriteReservations(dhcp.DefaultConfDir, d.vpcName+"_"+d.bridge, name, if err := writeDHCPFiles(d, name); err != nil {
[]dhcp.Reservation{{MAC: d.mac, IP: d.ip}}); err != nil { return err
return fmt.Errorf("write dhcp reservation: %w", err)
} }
if err := metadata.StartMetadata(metadata.NoCloudConfig{ if err := metadata.StartMetadata(metadata.NoCloudConfig{
Name: name, Name: name,
VpcName: d.vpcName, VpcName: nic.vpcName,
BindIP: d.interfaceIP, BindIP: nic.interfaceIP,
BindPort: d.metadataPort, BindPort: d.metadataPort,
Password: d.password, Password: d.password,
SSHKEY: d.sshkey, SSHKEY: d.sshkey,
@ -64,10 +75,14 @@ func StartVM(db *badger.DB, name string, cfg *configuration.Config) error {
qDisks[i] = qemu.DiskConfig{Path: disk.path, Dev: disk.dev} qDisks[i] = qemu.DiskConfig{Path: disk.path, Dev: disk.dev}
} }
qNICs := make([]qemu.NICConfig, len(d.nics))
for i, n := range d.nics {
qNICs[i] = qemu.NICConfig{TapID: n.tapID, Mac: n.mac}
}
qcfg := qemu.Config{ qcfg := qemu.Config{
Name: name, Name: name,
TapID: d.tapID, NICs: qNICs,
Mac: d.mac,
Disks: qDisks, Disks: qDisks,
Memory: d.memory, Memory: d.memory,
CPUs: d.cpus, CPUs: d.cpus,
@ -85,7 +100,7 @@ func StartVM(db *badger.DB, name string, cfg *configuration.Config) error {
qcfg.UEFIVarsPath = varsPath qcfg.UEFIVarsPath = varsPath
} }
if err := netns.Call(d.vpcName, func() error { if err := netns.Call(nic.vpcName, func() error {
return qemu.Start(qcfg) return qemu.Start(qcfg)
}); err != nil { }); err != nil {
return fmt.Errorf("start qemu: %w", err) return fmt.Errorf("start qemu: %w", err)
@ -94,6 +109,53 @@ func StartVM(db *badger.DB, name string, cfg *configuration.Config) error {
return state.Set(db, "vm/"+name, state.Running) return state.Set(db, "vm/"+name, state.Running)
} }
// writeDHCPFiles écrit, pour chaque subnet touché par la VM, les réservations
// de ses interfaces et les options qui suppriment la route par défaut sur les
// interfaces non primaires. Le subnet de l'interface primaire ne reçoit aucune
// option : les options non taggées du subnet portent déjà la route par défaut.
func writeDHCPFiles(d vmData, name string) error {
type subnetFiles struct {
nic nicData
reservations []dhcp.Reservation
tags []string
}
bySubnet := make(map[string]*subnetFiles)
for _, n := range d.nics {
confName := n.vpcName + "_" + n.bridge
if bySubnet[confName] == nil {
bySubnet[confName] = &subnetFiles{nic: n}
}
f := bySubnet[confName]
f.reservations = append(f.reservations, dhcp.Reservation{
MAC: n.mac, IP: n.ip, Tag: nicTag(name, n.index),
})
if !n.primary {
f.tags = append(f.tags, nicTag(name, n.index))
}
}
for confName, f := range bySubnet {
if err := dhcp.WriteReservations(dhcp.DefaultConfDir, confName, name, f.reservations); err != nil {
return fmt.Errorf("write dhcp reservations on %s: %w", confName, err)
}
if err := dhcp.WriteVMOptions(dhcp.DefaultConfDir, confName, name, f.tags, dhcp.Config{
InterfaceIP: net.ParseIP(f.nic.interfaceIP),
VPCRoute: f.nic.vpcCIDR,
}); err != nil {
return fmt.Errorf("write dhcp options on %s: %w", confName, err)
}
}
return nil
}
// nicTag identifie une interface auprès de dnsmasq. Il est par interface et non
// par VM : deux interfaces d'une même VM peuvent partager un subnet, et n'y
// avoir pas le même rôle.
func nicTag(vmName string, index int) string {
return fmt.Sprintf("%s-%d", vmName, index)
}
func copyFile(src, dst string) error { func copyFile(src, dst string) error {
if err := os.MkdirAll(filepath.Dir(dst), 0755); err != nil { if err := os.MkdirAll(filepath.Dir(dst), 0755); err != nil {
return err return err

View file

@ -3,6 +3,8 @@ package vm
import ( import (
"fmt" "fmt"
"math/rand" "math/rand"
"net"
"sort"
"strconv" "strconv"
"strings" "strings"
@ -16,15 +18,23 @@ type diskEntry struct {
dev string dev string
} }
type vmData struct { type nicData struct {
index int
subnetName string subnetName string
vpcName string vpcName string
interfaceIP string
bridge string bridge string
tapID int interfaceIP string
ip string ip string
metadataPort string
mac string mac string
tapID int
primary bool
mode string
vpcCIDR *net.IPNet
}
type vmData struct {
nics []nicData
metadataPort string
disks []diskEntry disks []diskEntry
memory int memory int
cpus int cpus int
@ -34,47 +44,25 @@ type vmData struct {
documents map[string]string documents map[string]string
} }
// primary retourne l'interface portant la route par défaut et le serveur de
// métadonnées. loadVM garantit qu'il y en a exactement une.
func (d vmData) primary() nicData {
for _, n := range d.nics {
if n.primary {
return n
}
}
return nicData{}
}
func loadVM(db *badger.DB, name string) (vmData, error) { func loadVM(db *badger.DB, name string) (vmData, error) {
var d vmData var d vmData
subnetName, err := kv.GetFromDB(db, "vm/"+name+"/subnet") nics, err := loadNICs(db, name)
if err != nil { if err != nil {
return d, fmt.Errorf("get subnet: %w", err) return d, err
} }
d.subnetName = subnetName d.nics = nics
d.bridge = "br-" + strings.SplitN(subnetName, "-", 2)[1]
vpcName, err := kv.GetFromDB(db, "subnet/"+subnetName+"/vpc")
if err != nil {
return d, fmt.Errorf("get vpc: %w", err)
}
d.vpcName = vpcName
interfaceIP, err := kv.GetFromDB(db, "subnet/"+subnetName+"/interface_ip")
if err != nil {
return d, fmt.Errorf("get interface_ip: %w", err)
}
d.interfaceIP = interfaceIP
tapIDStr, err := kv.GetFromDB(db, "vm/"+name+"/tap_id")
if err != nil {
d.tapID = rand.Intn(90000000) + 10000000
if err := kv.AddInDB(db, "vm/"+name+"/tap_id", strconv.Itoa(d.tapID)); err != nil {
return d, fmt.Errorf("store tap_id: %w", err)
}
} else {
tapID, err := strconv.Atoi(tapIDStr)
if err != nil {
return d, fmt.Errorf("parse tap_id: %w", err)
}
d.tapID = tapID
}
ip, err := kv.GetFromDB(db, "vm/"+name+"/ip")
if err != nil {
return d, fmt.Errorf("get ip: %w", err)
}
d.ip = ip
metadataPort, err := kv.GetFromDB(db, "vm/"+name+"/metadata_port") metadataPort, err := kv.GetFromDB(db, "vm/"+name+"/metadata_port")
if err != nil { if err != nil {
@ -82,12 +70,6 @@ func loadVM(db *badger.DB, name string) (vmData, error) {
} }
d.metadataPort = metadataPort d.metadataPort = metadataPort
mac, err := dhcp.GetMACForIP(db, d.subnetName, d.ip)
if err != nil {
return d, fmt.Errorf("get mac for ip %s: %w", d.ip, err)
}
d.mac = mac
diskEntries, err := kv.ListByPrefix(db, "vm/"+name+"/disk/") diskEntries, err := kv.ListByPrefix(db, "vm/"+name+"/disk/")
if err != nil { if err != nil {
return d, fmt.Errorf("list disks: %w", err) return d, fmt.Errorf("list disks: %w", err)
@ -139,3 +121,118 @@ func loadVM(db *badger.DB, name string) (vmData, error) {
return d, nil return d, nil
} }
// loadNICs lit les interfaces d'une VM sous vm/<name>/nic/<index>/.
// Le tap est alloué à la première lecture et persisté, comme avant le passage
// au multi-interfaces — mais désormais par interface.
func loadNICs(db *badger.DB, name string) ([]nicData, error) {
prefix := "vm/" + name + "/nic/"
entries, err := kv.ListByPrefix(db, prefix)
if err != nil {
return nil, fmt.Errorf("list nics: %w", err)
}
indexes := make(map[int]bool)
for key := range entries {
parts := strings.Split(strings.TrimPrefix(key, prefix), "/")
if len(parts) != 2 {
continue
}
idx, err := strconv.Atoi(parts[0])
if err != nil {
return nil, fmt.Errorf("invalid nic index %q for vm %s", parts[0], name)
}
indexes[idx] = true
}
if len(indexes) == 0 {
return nil, fmt.Errorf("no interface found for vm %q", name)
}
nics := make([]nicData, 0, len(indexes))
for idx := range indexes {
n, err := loadNIC(db, name, idx, entries)
if err != nil {
return nil, err
}
nics = append(nics, n)
}
sort.Slice(nics, func(i, j int) bool { return nics[i].index < nics[j].index })
primaries := 0
for _, n := range nics {
if n.primary {
primaries++
}
}
if primaries != 1 {
return nil, fmt.Errorf("vm %q has %d primary interfaces, expected exactly one", name, primaries)
}
return nics, nil
}
func loadNIC(db *badger.DB, name string, idx int, entries map[string]string) (nicData, error) {
n := nicData{index: idx}
prefix := fmt.Sprintf("vm/%s/nic/%d/", name, idx)
n.subnetName = entries[prefix+"subnet"]
if n.subnetName == "" {
return n, fmt.Errorf("nic %d of vm %s has no subnet", idx, name)
}
n.bridge = "br-" + strings.SplitN(n.subnetName, "-", 2)[1]
n.primary = entries[prefix+"primary"] == "true"
vpcName, err := kv.GetFromDB(db, "subnet/"+n.subnetName+"/vpc")
if err != nil {
return n, fmt.Errorf("get vpc of subnet %s: %w", n.subnetName, err)
}
n.vpcName = vpcName
interfaceIP, err := kv.GetFromDB(db, "subnet/"+n.subnetName+"/interface_ip")
if err != nil {
return n, fmt.Errorf("get interface_ip of subnet %s: %w", n.subnetName, err)
}
n.interfaceIP = interfaceIP
n.mode, err = kv.GetFromDB(db, "subnet/"+n.subnetName+"/mode")
if err != nil {
return n, fmt.Errorf("get mode of subnet %s: %w", n.subnetName, err)
}
if n.mode != "bridge" {
cidrStr, err := kv.GetFromDB(db, "vpc/"+n.vpcName+"/cidr")
if err != nil {
return n, fmt.Errorf("get cidr of vpc %s: %w", n.vpcName, err)
}
_, vpcCIDR, err := net.ParseCIDR(cidrStr)
if err != nil {
return n, fmt.Errorf("parse cidr of vpc %s: %w", n.vpcName, err)
}
n.vpcCIDR = vpcCIDR
}
n.ip = entries[prefix+"ip"]
if n.ip == "" {
return n, fmt.Errorf("nic %d of vm %s has no ip", idx, name)
}
mac, err := dhcp.GetMACForIP(db, n.subnetName, n.ip)
if err != nil {
return n, fmt.Errorf("get mac for ip %s: %w", n.ip, err)
}
n.mac = mac
if tapIDStr, ok := entries[prefix+"tap_id"]; ok {
tapID, err := strconv.Atoi(tapIDStr)
if err != nil {
return n, fmt.Errorf("parse tap_id of nic %d: %w", idx, err)
}
n.tapID = tapID
return n, nil
}
n.tapID = rand.Intn(90000000) + 10000000
if err := kv.AddInDB(db, prefix+"tap_id", strconv.Itoa(n.tapID)); err != nil {
return n, fmt.Errorf("store tap_id of nic %d: %w", idx, err)
}
return n, nil
}

View file

@ -1,6 +1,8 @@
package vm package vm
import ( import (
"fmt"
"strconv"
"testing" "testing"
"git.g3e.fr/syonad/two/pkg/db/kv" "git.g3e.fr/syonad/two/pkg/db/kv"
@ -12,11 +14,14 @@ func newVMInDB(t *testing.T) *badger.DB {
db := kv.InitDB(kv.Config{Path: t.TempDir()}, false) db := kv.InitDB(kv.Config{Path: t.TempDir()}, false)
t.Cleanup(func() { db.Close() }) t.Cleanup(func() { db.Close() })
kv.AddInDB(db, "vm/vm-1/subnet", "sn-000001") kv.AddInDB(db, "vm/vm-1/nic/0/subnet", "sn-000001")
kv.AddInDB(db, "vm/vm-1/nic/0/primary", "true")
kv.AddInDB(db, "subnet/sn-000001/vpc", "vp-admin") kv.AddInDB(db, "subnet/sn-000001/vpc", "vp-admin")
kv.AddInDB(db, "subnet/sn-000001/interface_ip", "10.1.1.1") kv.AddInDB(db, "subnet/sn-000001/interface_ip", "10.1.1.1")
kv.AddInDB(db, "subnet/sn-000001/mode", "vxlan")
kv.AddInDB(db, "vpc/vp-admin/cidr", "192.168.0.0/16")
kv.AddInDB(db, "subnet/sn-000001/dhcp/10.1.1.2", "00:22:33:00:01:02") kv.AddInDB(db, "subnet/sn-000001/dhcp/10.1.1.2", "00:22:33:00:01:02")
kv.AddInDB(db, "vm/vm-1/ip", "10.1.1.2") kv.AddInDB(db, "vm/vm-1/nic/0/ip", "10.1.1.2")
kv.AddInDB(db, "vm/vm-1/metadata_port", "8081") kv.AddInDB(db, "vm/vm-1/metadata_port", "8081")
kv.AddInDB(db, "vm/vm-1/disk/vda", "/data/root.qcow2") kv.AddInDB(db, "vm/vm-1/disk/vda", "/data/root.qcow2")
kv.AddInDB(db, "vm/vm-1/memory", "2048") kv.AddInDB(db, "vm/vm-1/memory", "2048")
@ -88,3 +93,109 @@ func TestLoadVM_EmptyDocumentIsPreserved(t *testing.T) {
t.Errorf("contenu attendu vide, obtenu %q", content) t.Errorf("contenu attendu vide, obtenu %q", content)
} }
} }
// --- multi-interfaces ---
func addNIC(t *testing.T, db *badger.DB, idx int, subnet, ip, mac string, primary bool) {
t.Helper()
prefix := fmt.Sprintf("vm/vm-1/nic/%d/", idx)
kv.AddInDB(db, prefix+"subnet", subnet)
kv.AddInDB(db, prefix+"ip", ip)
if primary {
kv.AddInDB(db, prefix+"primary", "true")
}
kv.AddInDB(db, "subnet/"+subnet+"/vpc", "vp-admin")
kv.AddInDB(db, "subnet/"+subnet+"/interface_ip", "10.1.1.1")
kv.AddInDB(db, "subnet/"+subnet+"/mode", "vxlan")
kv.AddInDB(db, "vpc/vp-admin/cidr", "192.168.0.0/16")
kv.AddInDB(db, "subnet/"+subnet+"/dhcp/"+ip, mac)
}
func TestLoadVM_SingleNIC(t *testing.T) {
d, err := loadVM(newVMInDB(t), "vm-1")
if err != nil {
t.Fatalf("loadVM : %v", err)
}
if len(d.nics) != 1 {
t.Fatalf("1 interface attendue, obtenu %d", len(d.nics))
}
if !d.primary().primary || d.primary().ip != "10.1.1.2" {
t.Errorf("primaire inattendue : %+v", d.primary())
}
}
func TestLoadVM_MultipleNICsSortedByIndex(t *testing.T) {
db := newVMInDB(t)
addNIC(t, db, 2, "sn-000003", "10.3.0.9", "00:22:33:00:03:09", false)
addNIC(t, db, 1, "sn-000002", "10.2.0.5", "00:22:33:00:02:05", false)
d, err := loadVM(db, "vm-1")
if err != nil {
t.Fatalf("loadVM : %v", err)
}
if len(d.nics) != 3 {
t.Fatalf("3 interfaces attendues, obtenu %d", len(d.nics))
}
for i, n := range d.nics {
if n.index != i {
t.Errorf("interface en position %d porte l'index %d — l'ordre détermine le slot PCI", i, n.index)
}
}
}
func TestLoadVM_TapAllocatedPerNIC(t *testing.T) {
db := newVMInDB(t)
addNIC(t, db, 1, "sn-000002", "10.2.0.5", "00:22:33:00:02:05", false)
d, err := loadVM(db, "vm-1")
if err != nil {
t.Fatalf("loadVM : %v", err)
}
if d.nics[0].tapID == d.nics[1].tapID {
t.Errorf("deux interfaces partagent le tap %d", d.nics[0].tapID)
}
for _, n := range d.nics {
stored, err := kv.GetFromDB(db, fmt.Sprintf("vm/vm-1/nic/%d/tap_id", n.index))
if err != nil {
t.Errorf("tap_id de l'interface %d non persisté : %v", n.index, err)
continue
}
if stored != strconv.Itoa(n.tapID) {
t.Errorf("tap_id de l'interface %d : %q en base, %d en mémoire", n.index, stored, n.tapID)
}
}
}
func TestLoadVM_NoPrimaryIsAnError(t *testing.T) {
db := kv.InitDB(kv.Config{Path: t.TempDir()}, false)
t.Cleanup(func() { db.Close() })
kv.AddInDB(db, "vm/vm-1/metadata_port", "8081")
kv.AddInDB(db, "vm/vm-1/disk/vda", "/data/root.qcow2")
kv.AddInDB(db, "vm/vm-1/memory", "2048")
kv.AddInDB(db, "vm/vm-1/cpus", "2")
addNIC(t, db, 0, "sn-000001", "10.1.1.2", "00:22:33:00:01:02", false)
kv.AddInDB(db, "subnet/sn-000001/mode", "vxlan")
if _, err := loadVM(db, "vm-1"); err == nil {
t.Error("aucune interface primaire : loadVM doit échouer plutôt que de laisser StartVM choisir au hasard")
}
}
func TestLoadVM_TwoPrimariesIsAnError(t *testing.T) {
db := newVMInDB(t)
addNIC(t, db, 1, "sn-000002", "10.2.0.5", "00:22:33:00:02:05", true)
if _, err := loadVM(db, "vm-1"); err == nil {
t.Error("deux interfaces primaires doivent être refusées")
}
}
func TestLoadVM_NoNICIsAnError(t *testing.T) {
db := kv.InitDB(kv.Config{Path: t.TempDir()}, false)
t.Cleanup(func() { db.Close() })
kv.AddInDB(db, "vm/vm-1/memory", "2048")
if _, err := loadVM(db, "vm-1"); err == nil {
t.Error("une VM sans interface doit être refusée")
}
}

View file

@ -32,6 +32,7 @@ func StopVM(db *badger.DB, name string, cfg *configuration.Config) error {
if err != nil { if err != nil {
return err return err
} }
nic := d.primary()
socketPath := filepath.Join(cfg.QEMU.QMPDir, name+".sock") socketPath := filepath.Join(cfg.QEMU.QMPDir, name+".sock")
@ -47,8 +48,13 @@ func StopVM(db *badger.DB, name string, cfg *configuration.Config) error {
} }
// socket absent ou QEMU déjà arrêté : cleanup direct // socket absent ou QEMU déjà arrêté : cleanup direct
if err := netns.Call(d.vpcName, func() error { if err := netns.Call(nic.vpcName, func() error {
return iptables.DeleteMetadataRedirect(d.ip, d.interfaceIP, d.metadataPort) for _, n := range d.nics {
if err := iptables.DeleteMetadataRedirect(n.ip, nic.interfaceIP, d.metadataPort); err != nil {
return fmt.Errorf("interface %d: %w", n.index, err)
}
}
return nil
}); err != nil { }); err != nil {
return fmt.Errorf("delete metadata redirect: %w", err) return fmt.Errorf("delete metadata redirect: %w", err)
} }
@ -57,11 +63,13 @@ func StopVM(db *badger.DB, name string, cfg *configuration.Config) error {
return fmt.Errorf("stop metadata: %w", err) return fmt.Errorf("stop metadata: %w", err)
} }
if err := netif.DeleteTap(d.tapID, d.vpcName); err != nil { for _, n := range d.nics {
return fmt.Errorf("delete tap: %w", err) if err := netif.DeleteTap(n.tapID, n.vpcName); err != nil {
return fmt.Errorf("delete tap of interface %d: %w", n.index, err)
}
} }
if err := removeDHCPReservation(d, name); err != nil { if err := removeDHCPFiles(d, name); err != nil {
return err return err
} }
@ -73,12 +81,26 @@ func StopVM(db *badger.DB, name string, cfg *configuration.Config) error {
return state.Set(db, "vm/"+name, state.Deleted) return state.Set(db, "vm/"+name, state.Deleted)
} }
// removeDHCPReservation retire le fichier de réservation puis redémarre dnsmasq : // removeDHCPFiles retire les fichiers de la VM dans chaque subnet qu'elle
// un fichier ajouté dans un dhcp-hostsdir est relu à chaud, un fichier retiré ne // touche, puis redémarre les dnsmasq concernés : un fichier ajouté dans un
// l'est pas (vérifié sur dnsmasq 2.90). // dhcp-hostsdir est relu à chaud, un fichier retiré ne l'est pas (vérifié sur
func removeDHCPReservation(d vmData, name string) error { // dnsmasq 2.90).
confName := d.vpcName + "_" + d.bridge func removeDHCPFiles(d vmData, name string) error {
seen := make(map[string]bool)
for _, n := range d.nics {
confName := n.vpcName + "_" + n.bridge
if seen[confName] {
continue
}
seen[confName] = true
if err := removeDHCPReservation(confName, name); err != nil {
return err
}
}
return nil
}
func removeDHCPReservation(confName, name string) error {
if err := dhcp.RemoveReservations(dhcp.DefaultConfDir, confName, name); err != nil { if err := dhcp.RemoveReservations(dhcp.DefaultConfDir, confName, name); err != nil {
return err return err
} }

View file

@ -4,7 +4,9 @@ import (
"errors" "errors"
"fmt" "fmt"
"path/filepath" "path/filepath"
"sort"
"strconv" "strconv"
"strings"
configuration "git.g3e.fr/syonad/two/internal/config/agent" configuration "git.g3e.fr/syonad/two/internal/config/agent"
"git.g3e.fr/syonad/two/internal/netns" "git.g3e.fr/syonad/two/internal/netns"
@ -46,33 +48,65 @@ func CheckVMs(db *badger.DB, cfg *configuration.Config, u unitChecker, n notify.
} }
func checkVM(db *badger.DB, cfg *configuration.Config, name string, u unitChecker, n notify.Notifier) { func checkVM(db *badger.DB, cfg *configuration.Config, name string, u unitChecker, n notify.Notifier) {
subnetName, err := kv.GetFromDB(db, prefixVM+name+"/subnet") entries, err := kv.ListByPrefix(db, prefixVM+name+"/nic/")
if err != nil { if err != nil {
n.Notify(kindVM, name, fmt.Sprintf("subnet unreadable in database: %v", err)) n.Notify(kindVM, name, fmt.Sprintf("interfaces unreadable in database: %v", err))
return
}
indexes := nicIndexes(entries, prefixVM+name+"/nic/")
if len(indexes) == 0 {
n.Notify(kindVM, name, "no interface in database")
return return
} }
for _, idx := range indexes {
prefix := fmt.Sprintf("%s%s/nic/%d/", prefixVM, name, idx)
subnetName := entries[prefix+"subnet"]
if subnetName == "" {
n.Notify(kindVM, name, fmt.Sprintf("interface %d has no subnet in database", idx))
continue
}
vpc, err := kv.GetFromDB(db, prefixSubnet+subnetName+"/vpc") vpc, err := kv.GetFromDB(db, prefixSubnet+subnetName+"/vpc")
if err != nil { if err != nil {
n.Notify(kindVM, name, fmt.Sprintf("vpc of subnet %s unreadable in database: %v", subnetName, err)) n.Notify(kindVM, name, fmt.Sprintf("vpc of subnet %s unreadable in database: %v", subnetName, err))
return continue
}
checkVMTap(entries[prefix+"tap_id"], name, vpc, idx, n)
} }
checkVMTap(db, name, vpc, n)
checkVMQemu(cfg, name, n) checkVMQemu(cfg, name, n)
checkUnit(kindVM, name, "metadata@"+name+".service", u, n) checkUnit(kindVM, name, "metadata@"+name+".service", u, n)
checkUnit(kindVM, name, qemu.ScopeName(name), u, n) checkUnit(kindVM, name, qemu.ScopeName(name), u, n)
} }
func checkVMTap(db *badger.DB, name, vpc string, n notify.Notifier) { // nicIndexes retourne les index d'interface présents en base, triés.
raw, err := kv.GetFromDB(db, prefixVM+name+"/tap_id") func nicIndexes(entries map[string]string, prefix string) []int {
if err != nil { seen := make(map[int]bool)
n.Notify(kindVM, name, fmt.Sprintf("tap_id unreadable in database: %v", err)) for key := range entries {
parts := strings.Split(strings.TrimPrefix(key, prefix), "/")
if len(parts) != 2 {
continue
}
if idx, err := strconv.Atoi(parts[0]); err == nil {
seen[idx] = true
}
}
indexes := make([]int, 0, len(seen))
for idx := range seen {
indexes = append(indexes, idx)
}
sort.Ints(indexes)
return indexes
}
func checkVMTap(raw, name, vpc string, idx int, n notify.Notifier) {
if raw == "" {
n.Notify(kindVM, name, fmt.Sprintf("interface %d has no tap_id in database", idx))
return return
} }
tapID, err := strconv.Atoi(raw) tapID, err := strconv.Atoi(raw)
if err != nil { if err != nil {
n.Notify(kindVM, name, fmt.Sprintf("invalid tap_id %q: %v", raw, err)) n.Notify(kindVM, name, fmt.Sprintf("interface %d has an invalid tap_id %q: %v", idx, raw, err))
return return
} }

View file

@ -20,10 +20,11 @@ func testCfg(t *testing.T) *configuration.Config {
func seedVM(t *testing.T, db *badger.DB, name, subnetName, vpc, tapID string) { func seedVM(t *testing.T, db *badger.DB, name, subnetName, vpc, tapID string) {
t.Helper() t.Helper()
seedResource(t, db, prefixVM, name, state.Running) seedResource(t, db, prefixVM, name, state.Running)
seedKV(t, db, prefixVM+name+"/subnet", subnetName) seedKV(t, db, prefixVM+name+"/nic/0/subnet", subnetName)
seedKV(t, db, prefixVM+name+"/nic/0/primary", "true")
seedKV(t, db, prefixSubnet+subnetName+"/vpc", vpc) seedKV(t, db, prefixSubnet+subnetName+"/vpc", vpc)
if tapID != "" { if tapID != "" {
seedKV(t, db, prefixVM+name+"/tap_id", tapID) seedKV(t, db, prefixVM+name+"/nic/0/tap_id", tapID)
} }
} }
@ -73,7 +74,7 @@ func TestCheckVMs_IgnoreLesEtatsNonRunning(t *testing.T) {
} }
} }
func TestCheckVMs_SubnetManquantEnBase(t *testing.T) { func TestCheckVMs_AucuneInterfaceEnBase(t *testing.T) {
db := newTestDB(t) db := newTestDB(t)
seedResource(t, db, prefixVM, "i-test1", state.Running) seedResource(t, db, prefixVM, "i-test1", state.Running)
r := &recorder{} r := &recorder{}
@ -86,15 +87,16 @@ func TestCheckVMs_SubnetManquantEnBase(t *testing.T) {
if len(got) != 1 { if len(got) != 1 {
t.Fatalf("attendu 1 notification, obtenu %d : %v", len(got), got) t.Fatalf("attendu 1 notification, obtenu %d : %v", len(got), got)
} }
if !strings.Contains(got[0].problem, "subnet unreadable") { if !strings.Contains(got[0].problem, "no interface") {
t.Errorf("problem = %q, devrait porter sur le subnet", got[0].problem) t.Errorf("problem = %q, devrait signaler l'absence d'interface", got[0].problem)
} }
} }
func TestCheckVMs_VPCDuSubnetManquant(t *testing.T) { func TestCheckVMs_VPCDuSubnetManquant(t *testing.T) {
db := newTestDB(t) db := newTestDB(t)
seedResource(t, db, prefixVM, "i-test1", state.Running) seedResource(t, db, prefixVM, "i-test1", state.Running)
seedKV(t, db, prefixVM+"i-test1/subnet", "br-000042") seedKV(t, db, prefixVM+"i-test1/nic/0/subnet", "br-000042")
seedKV(t, db, prefixVM+"i-test1/nic/0/primary", "true")
r := &recorder{} r := &recorder{}
if err := CheckVMs(db, testCfg(t), newFakeUnits(), r); err != nil { if err := CheckVMs(db, testCfg(t), newFakeUnits(), r); err != nil {
@ -115,8 +117,8 @@ func TestCheckVMs_TapIDManquant(t *testing.T) {
t.Fatalf("erreur inattendue: %v", err) t.Fatalf("erreur inattendue: %v", err)
} }
if !r.hasProblemContaining("tap_id unreadable") { if !r.hasProblemContaining("has no tap_id") {
t.Errorf("devrait signaler un tap_id illisible, obtenu %v", r.calls) t.Errorf("devrait signaler un tap_id absent, obtenu %v", r.calls)
} }
} }