Compare commits

...

8 commits

Author SHA1 Message Date
7138a6654a
Merge branch 'feature-33' 2026-08-25 23:57:35 +02:00
c47577db21
f-33: doc: update README
Signed-off-by: GnomeZworc <nicolas.boufidjeline@g3e.fr>
2026-08-25 23:57:11 +02:00
642a9c2f7a
f-33: doc: update release note
Signed-off-by: GnomeZworc <nicolas.boufidjeline@g3e.fr>
2026-08-25 23:56:55 +02:00
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
1ec1d44c1a
f-33: mise en place d'un nouveau system de reservation #33
Signed-off-by: GnomeZworc <nicolas.boufidjeline@g3e.fr>
2026-08-25 22:37:07 +02:00
f90514a5dc
f-33: dnsmasq: remove duplicated --no-daemon flag
Le drapeau était passé deux fois, en tête et en fin de ligne de commande.
Sans effet, mais gênant à la lecture.

Signed-off-by: GnomeZworc <nicolas.boufidjeline@g3e.fr>
2026-08-25 21:42:08 +02:00
30bb4107a1
f-33: dnsmasq: give each instance its own lease file
Sans --dhcp-leasefile, toutes les instances écrivaient dans le fichier de
baux par défaut du système : chacune relisait au démarrage des baux
appartenant aux autres subnets, puis réécrivait le fichier avec les siens
seulement.

Sans effet visible jusqu'ici — les réservations étant statiques,
l'adressage vient de la MAC et non du bail — mais la collision est réelle
et se manifeste à chaque redémarrage d'une instance.

Le fichier suit le même motif que le pidfile déjà en place, et vit dans
/run : ces baux n'ont aucun sens sans le netns, qui ne survit pas au
redémarrage de l'host.

Signed-off-by: GnomeZworc <nicolas.boufidjeline@g3e.fr>
2026-08-25 21:41:56 +02:00
26 changed files with 1456 additions and 193 deletions

View file

@ -1,19 +1,80 @@
# syonad # syonad/two
A simple but powerful orchestrator, designed to be easy to use and API-first. Orchestrateur réseau et machines virtuelles mono-nœud, pensé pour être piloté par un logiciel
plutôt que par un humain.
## deployer Il expose une API HTTP qui crée des **VPC** — isolés par network namespace —, des **subnets** —
en VXLAN ou attachés à un bridge existant — et des **VM** QEMU/KVM raccordées à ces subnets, avec
DHCP, routage et metadata cloud-init fournis automatiquement.
## Installation
```bash
curl -O https://git.g3e.fr/syonad/two/raw/branch/main/scripts/deploy.sh
bash ./deploy.sh -t 0.1.0 -i
``` ```
curl 'https://git.g3e.fr/syonad/two/raw/branch/main/scripts/deploy.sh' -O
bash deploy.sh
mv deploy.sh /opt/two/bin/
curl 'https://git.g3e.fr/syonad/two/raw/branch/main/systemd/agent.service' -o '/etc/systemd/system/agent.service' `deploy.sh` se met à jour lui-même depuis la branche, télécharge binaires, units systemd et
curl 'https://git.g3e.fr/syonad/two/raw/branch/main/systemd/dnsmasq@.service' -o '/etc/systemd/system/dnsmasq@.service' scripts depuis la release, et les vérifie contre le manifeste `SHA256SUMS`. Le drapeau `-i`
curl 'https://git.g3e.fr/syonad/two/raw/branch/main/systemd/metadata@.service' -o '/etc/systemd/system/metadata@.service' prépare l'host : paquets, module `br_netfilter`, `sysctl`, et bridges.
systemctl daemon-reload
modprobe br_netfilter Options utiles :
echo br_netfilter > /etc/modules-load.d/br_netfilter.conf
| Option | Effet |
|---|---|
| `-t <tag>` | déployer une release donnée |
| `-b <branche>` | déployer depuis une branche au lieu d'une release |
| `-i` | préparer l'host (paquets, noyau, réseau) |
| `-u <iface>` | interface physique d'uplink, `eno1` par défaut |
| `-B <bridge>` | bridge principal auquel l'uplink est rattaché |
| `-d` | dry-run : affiche les commandes sans les exécuter |
| `-V` | désactiver la vérification des sommes de contrôle |
Un déploiement relève les instances `dnsmasq@` et `metadata@` actives **avant** l'arrêt des
services, et les redémarre ensuite — c'est la seule façon de savoir lesquelles relancer.
## Configuration
Un seul fichier, `/etc/two/agent.yml`, partagé par les trois binaires. Voir
[`conf/agent/config.exemple.yml`](conf/agent/config.exemple.yml) pour l'ensemble des options :
chemins de la base et des sockets QEMU, pool de workers, correspondance des types d'interface vers
les bridges physiques, watchdog, API d'administration, journalisation.
## Prise en main
```bash
# Un VPC, avec son CIDR interne
curl -X POST http://127.0.0.1:8080/vpcs \
-d '{"name": "vp-admin", "cidr": "192.168.0.0/16"}'
# Un subnet en VXLAN dans ce VPC
curl -X POST http://127.0.0.1:8080/subnets \
-d '{"name": "sn-000001", "vpc": "vp-admin", "mode": "vxlan", "vxlan_id": 1,
"iface_type": "vms", "interface_ip": "10.1.1.1", "cidr": "10.1.0.0/23"}'
# Une VM, avec une clé SSH et un user-data cloud-init en base64
curl -X POST http://127.0.0.1:8080/vms \
-d '{"name": "i-web", "memory": 2048, "cpus": 2,
"metadata": {"sshkey": "ssh-ed25519 AAAA…",
"user_data": "'"$(base64 -w0 < user-data.yml)"'"},
"interfaces": [{"subnet": "sn-000001", "ip": "10.1.1.2", "primary": true}],
"storage": [{"path": "/data/disks/vms/i-web.qcow2", "dev": "vda"}]}'
``` ```
Les créations sont **asynchrones** : l'API répond `202` et l'état de la ressource passe de
`creating` à `running` en base. `GET /vms/i-web` renvoie l'état courant.
Une VM peut porter plusieurs interfaces, dans un même VPC ; exactement une doit être marquée
`primary` — elle porte la route par défaut et le serveur de metadata.
La spécification complète est dans [`api/agent.yaml`](api/agent.yaml).
## Composants
| Binaire | Rôle |
|---|---|
| `agent` | processus principal : API, dispatcher, exécution, watchdog |
| `metadata` | serveur de metadata cloud-init, une instance par VM dans le netns du VPC |
| `db` | inspection de la base clé-valeur en ligne de commande |
L'agent prend `-config`, les deux autres `-conf`.

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

@ -171,32 +171,56 @@ func TestGenerateConfig_ContainsDhcpRange(t *testing.T) {
} }
} }
func TestGenerateConfig_OneHostEntryPerIP(t *testing.T) { func TestGenerateConfig_OneEntryPerIP(t *testing.T) {
// /29 = réseau + broadcast + 6 hôtes → 8 adresses // /29 = 8 adresses. Les entrées ne vont plus dans le fichier dnsmasq mais
conf := newConf(t, "10.0.0.0/29") // dans la map retournée, que StoreDHCPEntries écrit en base pour GetMACForIP.
path, _, _ := GenerateConfig(conf) _, entries, err := GenerateConfig(newConf(t, "10.0.0.0/29"))
content, _ := os.ReadFile(path) if err != nil {
t.Fatalf("GenerateConfig : %v", err)
}
if len(entries) != 8 {
t.Errorf("attendu 8 entrées ip->mac, obtenu %d", len(entries))
}
}
lines := strings.Split(string(content), "\n") func TestGenerateConfig_NoPreGeneratedHosts(t *testing.T) {
count := 0 // Une entrée dhcp-host pré-générée fait rejeter celle du dhcp-hostsdir
for _, l := range lines { // (« duplicate dhcp-host IP address »), sans erreur : la VM reçoit alors
if strings.HasPrefix(l, "dhcp-host=") { // les options non taggées. Vérifié sur dnsmasq 2.90.
count++ content := confLines(t, newConf(t, "10.0.0.0/29"))
if strings.Contains(content, "dhcp-host=") {
t.Errorf("aucune entrée dhcp-host ne doit être pré-générée :\n%s", content)
}
}
func TestGenerateConfig_PointsToDirs(t *testing.T) {
conf := newConf(t, "10.0.0.0/29")
content := confLines(t, conf)
for _, want := range []string{
"dhcp-hostsdir=" + HostsDir(conf.ConfDir, conf.Name),
"dhcp-optsdir=" + OptsDir(conf.ConfDir, conf.Name),
} {
if !strings.Contains(content, want) {
t.Errorf("%q absent :\n%s", want, content)
} }
} }
// /29 contient 8 adresses (0 à 7) for _, dir := range []string{HostsDir(conf.ConfDir, conf.Name), OptsDir(conf.ConfDir, conf.Name)} {
if count != 8 { if fi, err := os.Stat(dir); err != nil || !fi.IsDir() {
t.Errorf("attendu 8 entrées dhcp-host, obtenu %d", count) t.Errorf("répertoire %q non créé : %v", dir, err)
}
} }
} }
func TestGenerateConfig_MACPrefix(t *testing.T) { func TestGenerateConfig_MACPrefix(t *testing.T) {
conf := newConf(t, "10.0.0.0/30") // 4 adresses _, entries, err := GenerateConfig(newConf(t, "10.0.0.0/30"))
path, _, _ := GenerateConfig(conf) if err != nil {
content, _ := os.ReadFile(path) t.Fatalf("GenerateConfig : %v", err)
}
if !strings.Contains(string(content), "00:22:33:") { for ip, mac := range entries {
t.Errorf("préfixe MAC 00:22:33: absent :\n%s", content) if !strings.HasPrefix(mac, "00:22:33:") {
t.Errorf("mac de %s sans le préfixe 00:22:33: : %s", ip, mac)
}
} }
} }

View file

@ -25,21 +25,24 @@ func GenerateConfig(c Config) (string, map[string]string, error) {
} else { } else {
fmt.Fprintf(&sb, "dhcp-option=3\n") fmt.Fprintf(&sb, "dhcp-option=3\n")
} }
fmt.Fprintf(&sb, "dhcp-option=6,1.1.1.1,8.8.8.8\n\n") fmt.Fprintf(&sb, "dhcp-option=6,1.1.1.1,8.8.8.8\n")
fmt.Fprintf(&sb, "dhcp-hostsdir=%s\n", HostsDir(c.ConfDir, c.Name))
fmt.Fprintf(&sb, "dhcp-optsdir=%s\n", OptsDir(c.ConfDir, c.Name))
entries := make(map[string]string) entries := make(map[string]string)
i := 0 i := 0
for ip := cloneIP(c.Network.IP); c.Network.Contains(ip); incrementIP(ip) { for ip := cloneIP(c.Network.IP); c.Network.Contains(ip); incrementIP(ip) {
mac := fmt.Sprintf("00:22:33:%02X:%02X:%02X", (i>>16)&0xFF, (i>>8)&0xFF, i&0xFF) entries[ip.String()] = fmt.Sprintf("00:22:33:%02X:%02X:%02X", (i>>16)&0xFF, (i>>8)&0xFF, i&0xFF)
fmt.Fprintf(&sb, "dhcp-host=%s,%s\n", mac, ip)
entries[ip.String()] = mac
i++ i++
} }
outPath := filepath.Join(c.ConfDir, c.Name+".conf") for _, dir := range []string{c.ConfDir, HostsDir(c.ConfDir, c.Name), OptsDir(c.ConfDir, c.Name)} {
if err := os.MkdirAll(c.ConfDir, 0755); err != nil { if err := os.MkdirAll(dir, 0755); err != nil {
return "", nil, err return "", nil, fmt.Errorf("create %s: %w", dir, err)
} }
}
outPath := filepath.Join(c.ConfDir, c.Name+".conf")
return outPath, entries, os.WriteFile(outPath, []byte(sb.String()), 0644) return outPath, entries, os.WriteFile(outPath, []byte(sb.String()), 0644)
} }

View file

@ -0,0 +1,119 @@
package dhcp
import (
"fmt"
"os"
"path/filepath"
"strings"
)
type Reservation struct {
MAC string
IP string
Tag string // pose set:<Tag> sur l'entrée, pour cibler les options par interface
}
func HostsDir(confDir, name string) string {
return filepath.Join(confDir, name+".hosts.d")
}
func OptsDir(confDir, name string) string {
return filepath.Join(confDir, name+".opts.d")
}
func UnitName(name string) string {
return "dnsmasq@" + name + ".service"
}
func WriteReservations(confDir, name, vmName string, res []Reservation) error {
if len(res) == 0 {
return fmt.Errorf("no reservation for vm %q: it would get no address", vmName)
}
dir := HostsDir(confDir, name)
if err := os.MkdirAll(dir, 0755); err != nil {
return fmt.Errorf("create %s: %w", dir, err)
}
var sb strings.Builder
for _, r := range res {
if 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)
continue
}
fmt.Fprintf(&sb, "%s,%s,set:%s\n", r.MAC, r.IP, r.Tag)
}
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
}
// 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 {
for _, dir := range []string{HostsDir(confDir, name), OptsDir(confDir, name)} {
path := filepath.Join(dir, vmName)
if err := os.Remove(path); err != nil && !os.IsNotExist(err) {
return fmt.Errorf("remove %s: %w", path, err)
}
}
return nil
}
func RemoveSubnetDirs(confDir, name string) error {
for _, dir := range []string{HostsDir(confDir, name), OptsDir(confDir, name)} {
if err := os.RemoveAll(dir); err != nil {
return fmt.Errorf("remove %s: %w", dir, err)
}
}
return nil
}

View file

@ -0,0 +1,228 @@
package dhcp
import (
"net"
"os"
"path/filepath"
"strings"
"testing"
)
const subName = "vp-admin_br-000001"
func TestWriteReservations_WritesOneLinePerInterface(t *testing.T) {
dir := t.TempDir()
res := []Reservation{
{MAC: "00:22:33:00:01:02", IP: "10.1.1.2"},
{MAC: "00:22:33:00:02:07", IP: "10.1.2.7"},
}
if err := WriteReservations(dir, subName, "i-web", res); err != nil {
t.Fatalf("WriteReservations : %v", err)
}
content, err := os.ReadFile(filepath.Join(HostsDir(dir, subName), "i-web"))
if err != nil {
t.Fatalf("fichier absent : %v", err)
}
want := "00:22:33:00:01:02,10.1.1.2\n00:22:33:00:02:07,10.1.2.7\n"
if string(content) != want {
t.Errorf("contenu attendu %q, obtenu %q", want, content)
}
}
func TestWriteReservations_CreatesDir(t *testing.T) {
dir := t.TempDir()
if err := WriteReservations(dir, subName, "i-web", []Reservation{{MAC: "aa", IP: "10.0.0.1"}}); err != nil {
t.Fatalf("WriteReservations : %v", err)
}
if fi, err := os.Stat(HostsDir(dir, subName)); err != nil || !fi.IsDir() {
t.Errorf("hosts.d non créé : %v", err)
}
}
func TestWriteReservations_EmptyIsAnError(t *testing.T) {
if err := WriteReservations(t.TempDir(), subName, "i-web", nil); err == nil {
t.Error("sans réservation la VM n'obtiendrait aucune adresse : il faut échouer, pas écrire un fichier vide")
}
}
func TestWriteReservations_IncompleteIsAnError(t *testing.T) {
cases := []Reservation{{MAC: "", IP: "10.0.0.1"}, {MAC: "aa:bb", IP: ""}}
for _, r := range cases {
if err := WriteReservations(t.TempDir(), subName, "i-web", []Reservation{r}); err == nil {
t.Errorf("réservation incomplète acceptée : %+v", r)
}
}
}
func TestWriteReservations_Overwrites(t *testing.T) {
dir := t.TempDir()
_ = WriteReservations(dir, subName, "i-web", []Reservation{{MAC: "aa", IP: "10.0.0.1"}})
if err := WriteReservations(dir, subName, "i-web", []Reservation{{MAC: "bb", IP: "10.0.0.2"}}); err != nil {
t.Fatalf("WriteReservations : %v", err)
}
content, _ := os.ReadFile(filepath.Join(HostsDir(dir, subName), "i-web"))
if string(content) != "bb,10.0.0.2\n" {
t.Errorf("la réécriture doit remplacer, obtenu %q", content)
}
}
func TestRemoveReservations_RemovesBothFiles(t *testing.T) {
dir := t.TempDir()
_ = WriteReservations(dir, subName, "i-web", []Reservation{{MAC: "aa", IP: "10.0.0.1"}})
if err := os.MkdirAll(OptsDir(dir, subName), 0755); err != nil {
t.Fatal(err)
}
optsFile := filepath.Join(OptsDir(dir, subName), "i-web")
if err := os.WriteFile(optsFile, []byte("tag:i-web,3,10.0.0.1\n"), 0644); err != nil {
t.Fatal(err)
}
if err := RemoveReservations(dir, subName, "i-web"); err != nil {
t.Fatalf("RemoveReservations : %v", err)
}
for _, p := range []string{filepath.Join(HostsDir(dir, subName), "i-web"), optsFile} {
if _, err := os.Stat(p); !os.IsNotExist(err) {
t.Errorf("%s aurait dû être supprimé", p)
}
}
}
func TestRemoveReservations_AbsentIsNotAnError(t *testing.T) {
if err := RemoveReservations(t.TempDir(), subName, "jamais-creee"); err != nil {
t.Errorf("supprimer une VM sans réservation ne doit pas échouer : %v", err)
}
}
func TestRemoveReservations_LeavesOtherVMs(t *testing.T) {
dir := t.TempDir()
_ = WriteReservations(dir, subName, "i-web", []Reservation{{MAC: "aa", IP: "10.0.0.1"}})
_ = WriteReservations(dir, subName, "i-db", []Reservation{{MAC: "bb", IP: "10.0.0.2"}})
if err := RemoveReservations(dir, subName, "i-web"); err != nil {
t.Fatalf("RemoveReservations : %v", err)
}
if _, err := os.Stat(filepath.Join(HostsDir(dir, subName), "i-db")); err != nil {
t.Errorf("i-db ne devait pas être touchée : %v", err)
}
}
func TestRemoveSubnetDirs(t *testing.T) {
dir := t.TempDir()
_ = WriteReservations(dir, subName, "i-web", []Reservation{{MAC: "aa", IP: "10.0.0.1"}})
if err := RemoveSubnetDirs(dir, subName); err != nil {
t.Fatalf("RemoveSubnetDirs : %v", err)
}
for _, d := range []string{HostsDir(dir, subName), OptsDir(dir, subName)} {
if _, err := os.Stat(d); !os.IsNotExist(err) {
t.Errorf("%s aurait dû être supprimé", d)
}
}
}
func TestUnitName(t *testing.T) {
if got := UnitName(subName); got != "dnsmasq@vp-admin_br-000001.service" {
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

@ -68,6 +68,10 @@ func stopDHCP(db *badger.DB, subnetName string, d subnetData) error {
return fmt.Errorf("remove dnsmasq config: %w", err) return fmt.Errorf("remove dnsmasq config: %w", err)
} }
if err := dhcp.RemoveSubnetDirs(dhcp.DefaultConfDir, d.vpc+"_"+d.bridge); err != nil {
return fmt.Errorf("remove dnsmasq dirs: %w", err)
}
if err := kv.DeleteInDB(db, "subnet/"+subnetName+"/dhcp"); err != nil { if err := kv.DeleteInDB(db, "subnet/"+subnetName+"/dhcp"); err != nil {
return fmt.Errorf("delete dhcp entries: %w", err) return fmt.Errorf("delete dhcp entries: %w", err)
} }

View file

@ -3,10 +3,12 @@ package vm
import ( import (
"fmt" "fmt"
"io" "io"
"net"
"os" "os"
"path/filepath" "path/filepath"
configuration "git.g3e.fr/syonad/two/internal/config/agent" configuration "git.g3e.fr/syonad/two/internal/config/agent"
"git.g3e.fr/syonad/two/internal/dhcp"
"git.g3e.fr/syonad/two/internal/iptables" "git.g3e.fr/syonad/two/internal/iptables"
"git.g3e.fr/syonad/two/internal/metadata" "git.g3e.fr/syonad/two/internal/metadata"
"git.g3e.fr/syonad/two/internal/netif" "git.g3e.fr/syonad/two/internal/netif"
@ -30,21 +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 := writeDHCPFiles(d, name); err != nil {
return 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,
@ -58,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,
@ -79,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)
@ -88,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

@ -7,12 +7,14 @@ import (
"time" "time"
configuration "git.g3e.fr/syonad/two/internal/config/agent" configuration "git.g3e.fr/syonad/two/internal/config/agent"
"git.g3e.fr/syonad/two/internal/dhcp"
"git.g3e.fr/syonad/two/internal/iptables" "git.g3e.fr/syonad/two/internal/iptables"
"git.g3e.fr/syonad/two/internal/metadata" "git.g3e.fr/syonad/two/internal/metadata"
"git.g3e.fr/syonad/two/internal/netif" "git.g3e.fr/syonad/two/internal/netif"
"git.g3e.fr/syonad/two/internal/netns" "git.g3e.fr/syonad/two/internal/netns"
"git.g3e.fr/syonad/two/internal/qmp" "git.g3e.fr/syonad/two/internal/qmp"
"git.g3e.fr/syonad/two/internal/state" "git.g3e.fr/syonad/two/internal/state"
"git.g3e.fr/syonad/two/pkg/systemd"
"github.com/dgraph-io/badger/v4" "github.com/dgraph-io/badger/v4"
) )
@ -30,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")
@ -45,9 +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(nic.vpcName, func() error {
if err := netns.Call(d.vpcName, func() error { for _, n := range d.nics {
return iptables.DeleteMetadataRedirect(d.ip, d.interfaceIP, d.metadataPort) 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)
} }
@ -56,8 +63,14 @@ 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 := removeDHCPFiles(d, name); err != nil {
return err
} }
if d.uefi { if d.uefi {
@ -68,6 +81,52 @@ 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)
} }
// removeDHCPFiles retire les fichiers de la VM dans chaque subnet qu'elle
// touche, puis redémarre les dnsmasq concernés : un fichier ajouté dans un
// dhcp-hostsdir est relu à chaud, un fichier retiré ne l'est pas (vérifié sur
// dnsmasq 2.90).
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 {
return err
}
svc, err := systemd.New()
if err != nil {
return fmt.Errorf("connect to systemd: %w", err)
}
defer svc.Close()
unit := dhcp.UnitName(confName)
status, err := svc.Status(unit)
if err != nil || status.ActiveState != "active" {
return nil
}
if err := svc.Restart(unit); err != nil {
return fmt.Errorf("restart %s: %w", unit, err)
}
if status, err := svc.Status(unit); err != nil {
return fmt.Errorf("status %s after restart: %w", unit, err)
} else if status.ActiveState != "active" {
return fmt.Errorf("%s is %s after restart", unit, status.ActiveState)
}
return nil
}
func waitQMPDead(socketPath string, timeout, poll time.Duration) { func waitQMPDead(socketPath string, timeout, poll time.Duration) {
timer := time.After(timeout) timer := time.After(timeout)
for { for {

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

View file

@ -54,6 +54,11 @@ func (m *Manager) Stop(service string) error {
return m.job("StopUnit", service) return m.job("StopUnit", service)
} }
// Restart redémarre un service systemd
func (m *Manager) Restart(service string) error {
return m.job("RestartUnit", service)
}
func (m *Manager) job(method, service string) error { func (m *Manager) job(method, service string) error {
callCtx, callCancel := context.WithTimeout(context.Background(), defaultTimeout) callCtx, callCancel := context.WithTimeout(context.Background(), defaultTimeout)
defer callCancel() defer callCancel()
@ -66,6 +71,8 @@ func (m *Manager) job(method, service string) error {
_, err = m.conn.StartUnitContext(callCtx, service, jobMode, ch) _, err = m.conn.StartUnitContext(callCtx, service, jobMode, ch)
case "StopUnit": case "StopUnit":
_, err = m.conn.StopUnitContext(callCtx, service, jobMode, ch) _, err = m.conn.StopUnitContext(callCtx, service, jobMode, ch)
case "RestartUnit":
_, err = m.conn.RestartUnitContext(callCtx, service, jobMode, ch)
default: default:
return errors.New("unsupported job method") return errors.New("unsupported job method")
} }

View file

@ -11,23 +11,48 @@ Première version stable de **syonad/two**, orchestrateur réseau et VM mono-nœ
avec `error` en cas d'échec d'exécution ; suppression autorisée depuis `running` et `error` avec `error` en cas d'échec d'exécution ; suppression autorisée depuis `running` et `error`
- Migration au démarrage de l'agent : toute ressource restée dans un état transitoire est - Migration au démarrage de l'agent : toute ressource restée dans un état transitoire est
basculée en `error`, la file de travail étant en mémoire basculée en `error`, la file de travail étant en mémoire
- Arrêt gracieux : serveurs HTTP, puis drainage des workers, puis fermeture de la base. Si le
budget d'arrêt est dépassé, la base n'est **pas** fermée — la rejouer au démarrage suivant vaut
mieux que de la fermer sous un écrivain concurrent
**Réseau** **Réseau**
- VPC isolés par network namespace, subnets en mode `vxlan` ou `bridge` - VPC isolés par network namespace, subnets en mode `vxlan` ou `bridge`
- DHCP par subnet via instances `dnsmasq@` dédiées, entrées ip→mac en base - DHCP par subnet via instances `dnsmasq@` dédiées, entrées ip→mac en base, fichier de baux propre
- Route par défaut et route du VPC distribuées par DHCP (`default_route` par subnet) à chaque instance
- Toutes les routes sont distribuées par l'option 121 : route vers le serveur de metadata, route du
VPC, et route par défaut. L'option 3 reste émise pour les clients qui n'implémentent pas la 121
- `default_route` choisit le **next-hop** de la route par défaut : l'`interface_ip` du subnet par
défaut, sinon la `gateway` fournie ou celle déduite de l'host
- `gateway` optionnelle par subnet, non validée par l'agent
- Isolation du DHCP par ebtables, redirection du service de metadata par iptables - Isolation du DHCP par ebtables, redirection du service de metadata par iptables
**Machines virtuelles** **Machines virtuelles**
- Plusieurs interfaces réseau par VM, dans un même VPC. La position de l'interface détermine son
slot PCI, donc son nom dans le guest ; **exactement une** interface est primaire et porte la route
par défaut et le serveur de metadata
- Démarrage QEMU/KVM avec plusieurs disques et ordre de démarrage explicite - Démarrage QEMU/KVM avec plusieurs disques et ordre de démarrage explicite
- Amorçage UEFI optionnel (OVMF), avec magasin de variables par VM - Amorçage UEFI optionnel (OVMF), avec magasin de variables par VM
- Serveur de metadata cloud-init par VM (`metadata@`), sans base de données dans le processus - Serveur de metadata cloud-init par VM (`metadata@`), sans base de données dans le processus
- Les VMs survivent à l'arrêt de l'agent : QEMU est lancé hors de son cgroup via `systemd-run` - Les VMs survivent à l'arrêt de l'agent : QEMU est lancé hors de son cgroup via `systemd-run`
**Metadata cloud-init**
- Objet `metadata` dans la création de VM : `password`, `sshkey`, `user_data`
- `user_data` transmis en base64, ce qui autorise les charges gzip+base64 ; un encodage invalide est
refusé en 400, jamais servi vide en silence
- Un document fourni est servi **verbatim**, un document absent retombe sur le modèle par défaut, et
un document explicitement vide est servi vide — les trois cas sont distincts
- Le compte `syonad` n'est créé que si un mot de passe ou une clé est fourni, et reste verrouillé
quand seule une clé l'est. L'agent n'impose aucune modification du compte root
**Exploitation** **Exploitation**
- Watchdog de cohérence : vérifie périodiquement que les ressources `running` existent encore sur
le système et signale les écarts. Lecture seule, il ne répare jamais. Désactivé par défaut,
activé dans le fichier d'exemple
- API d'administration en lecture seule sur la boucle locale, pour inspecter la base
- Métriques Prometheus : nombre de VPC, subnets et VMs par état - Métriques Prometheus : nombre de VPC, subnets et VMs par état
- `deploy.sh` avec profils d'host (`kvm`), préparation système déléguée à `bootstrap_kvm.sh` - `deploy.sh` avec profils d'host (`kvm`), préparation système déléguée à `bootstrap_kvm.sh`
- Units systemd et scripts publiés comme assets de release, avec manifeste `SHA256SUMS` - Units systemd et scripts publiés comme assets de release, avec manifeste `SHA256SUMS`
@ -39,7 +64,17 @@ Première version stable de **syonad/two**, orchestrateur réseau et VM mono-nœ
- Pas de rollback en cas d'échec partiel d'une création — les ressources réseau orphelines - Pas de rollback en cas d'échec partiel d'une création — les ressources réseau orphelines
ne sont pas nettoyées automatiquement ne sont pas nettoyées automatiquement
- API destinée à un appelant logiciel : la validation de cohérence des entrées (CIDR, VXLAN - API destinée à un appelant logiciel : la validation de cohérence des entrées (CIDR, VXLAN
ID, format des noms) est à la charge de l'appelant ID, format des noms, joignabilité d'une `gateway`) est à la charge de l'appelant
- Le mode de subnet `public_ip` est accepté par l'API et par la sélection des routes DHCP, mais sa
mise en place réseau n'existe pas : créer un tel subnet échoue explicitement
- Les interfaces multiples d'une VM doivent appartenir au même VPC
- Le réseau des guests est configuré par le DHCP seul ; le `network-config` cloud-init servi est
sans effet et ne doit pas être « corrigé » sans mesurer l'impact sur les VM existantes
- Modifier les routes d'une VM déjà démarrée ne prend effet qu'au renouvellement du bail, soit
jusqu'à six heures plus tard, ou à son redémarrage
- `vm/<name>/password` contient un **hash**, stocké en clair en base et restitué par l'API
d'administration ; celle-ci est désactivée par défaut et n'écoute que sur la boucle locale
- L'API de l'agent n'a pas d'authentification : son exposition réseau doit être restreinte
- Les packages `internal/netns`, `netif`, `qemu`, `vm`, `iptables` et `ebtables` ne - Les packages `internal/netns`, `netif`, `qemu`, `vm`, `iptables` et `ebtables` ne
fonctionnent que sous Linux fonctionnent que sous Linux

View file

@ -10,10 +10,10 @@ echo "start dnsmasq ${NETNS} ${BRIDGE}"
exec ip netns exec "${NETNS}" \ exec ip netns exec "${NETNS}" \
dnsmasq \ dnsmasq \
--no-daemon \
--interface="${BRIDGE}" \ --interface="${BRIDGE}" \
--bind-interfaces \ --bind-interfaces \
--pid-file="/run/dnsmasq-$arg.pid" \ --pid-file="/run/dnsmasq-$arg.pid" \
--dhcp-leasefile="/run/dnsmasq-$arg.leases" \
--conf-file="/etc/dnsmasq.d/$arg.conf" \ --conf-file="/etc/dnsmasq.d/$arg.conf" \
--no-hosts \ --no-hosts \
--no-resolv \ --no-resolv \