diff --git a/internal/api/agent/vm.go b/internal/api/agent/vm.go index d9f6bf5..f932766 100644 --- a/internal/api/agent/vm.go +++ b/internal/api/agent/vm.go @@ -87,11 +87,8 @@ func vmFromDB(name string, entries map[string]string) (VM, error) { vm.Interfaces = []VMInterface{{Subnet: subnet, IP: ip, Primary: true}} } - diskPrefix := prefix + "disk/" - for key, path := range entries { - if dev := strings.TrimPrefix(key, diskPrefix); dev != key { - vm.Storage = append(vm.Storage, VMStorage{Path: path, Dev: dev}) - } + if path := entries[prefix+"volume_path"]; path != "" { + vm.Storage = []VMStorage{{Path: path}} } return vm, nil diff --git a/internal/api/agent/vm_test.go b/internal/api/agent/vm_test.go deleted file mode 100644 index 3af5968..0000000 --- a/internal/api/agent/vm_test.go +++ /dev/null @@ -1,158 +0,0 @@ -package agentapi - -import ( - "bytes" - "encoding/json" - "net/http" - "net/http/httptest" - "sort" - "testing" - - "git.g3e.fr/syonad/two/pkg/db/kv" -) - -// --- vmFromDB --- - -func TestVmFromDB_SingleDisk(t *testing.T) { - entries := map[string]string{ - "vm/vm-1/state": "started", - "vm/vm-1/subnet": "sn-1", - "vm/vm-1/ip": "10.0.0.5", - "vm/vm-1/metadata_port": "1234", - "vm/vm-1/memory": "512", - "vm/vm-1/cpus": "1", - "vm/vm-1/disk/sda": "/data/root.qcow2", - } - vm, err := vmFromDB("vm-1", entries) - if err != nil { - t.Fatalf("vmFromDB a échoué : %v", err) - } - if len(vm.Storage) != 1 { - t.Fatalf("attendu 1 disque, obtenu %d", len(vm.Storage)) - } - if vm.Storage[0].Dev != "sda" || vm.Storage[0].Path != "/data/root.qcow2" { - t.Errorf("disque inattendu : %+v", vm.Storage[0]) - } -} - -func TestVmFromDB_MultiDisk(t *testing.T) { - entries := map[string]string{ - "vm/vm-2/state": "started", - "vm/vm-2/subnet": "sn-1", - "vm/vm-2/ip": "10.0.0.6", - "vm/vm-2/metadata_port": "1235", - "vm/vm-2/memory": "1024", - "vm/vm-2/cpus": "2", - "vm/vm-2/disk/sda": "/data/root.qcow2", - "vm/vm-2/disk/sdb": "/data/data.qcow2", - } - vm, err := vmFromDB("vm-2", entries) - if err != nil { - t.Fatalf("vmFromDB a échoué : %v", err) - } - if len(vm.Storage) != 2 { - t.Fatalf("attendu 2 disques, obtenu %d", len(vm.Storage)) - } - sort.Slice(vm.Storage, func(i, j int) bool { return vm.Storage[i].Dev < vm.Storage[j].Dev }) - if vm.Storage[0].Dev != "sda" || vm.Storage[1].Dev != "sdb" { - t.Errorf("devs attendus [sda sdb], obtenus [%s %s]", vm.Storage[0].Dev, vm.Storage[1].Dev) - } -} - -func TestVmFromDB_SlotGap(t *testing.T) { - // sdb absent — sda et sdc seulement - entries := map[string]string{ - "vm/vm-3/state": "started", - "vm/vm-3/subnet": "sn-1", - "vm/vm-3/ip": "10.0.0.7", - "vm/vm-3/metadata_port": "1236", - "vm/vm-3/memory": "512", - "vm/vm-3/cpus": "1", - "vm/vm-3/disk/sda": "/data/root.qcow2", - "vm/vm-3/disk/sdc": "/data/extra.qcow2", - } - vm, err := vmFromDB("vm-3", entries) - if err != nil { - t.Fatalf("vmFromDB a échoué : %v", err) - } - if len(vm.Storage) != 2 { - t.Fatalf("attendu 2 disques, obtenu %d", len(vm.Storage)) - } - devs := map[string]bool{} - for _, s := range vm.Storage { - devs[s.Dev] = true - } - if !devs["sda"] || !devs["sdc"] { - t.Errorf("attendu sda et sdc, obtenus %v", devs) - } - if devs["sdb"] { - t.Error("sdb ne devrait pas apparaître") - } -} - -// --- POST /vms --- - -func TestStartVM_MultiDisk(t *testing.T) { - s, db := newTestServer(t) - kv.AddInDB(db, "subnet/sn-1/state", "created") - kv.AddInDB(db, "subnet/sn-1/vpc", "vpc-1") - - body, _ := json.Marshal(VMCreateRequest{ - Name: "vm-10", - Interfaces: []VMInterface{ - {Subnet: "sn-1", IP: "10.0.0.10", Primary: true}, - }, - Storage: []VMStorage{ - {Path: "/data/root.qcow2", Dev: "sda"}, - {Path: "/data/data.qcow2", Dev: "sdb"}, - }, - Memory: 1024, - CPUs: 2, - }) - - 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()) - } - - for _, dev := range []string{"sda", "sdb"} { - if _, err := kv.GetFromDB(db, "vm/vm-10/disk/"+dev); err != nil { - t.Errorf("disk/%s absent en DB après création", dev) - } - } - if _, err := kv.GetFromDB(db, "vm/vm-10/volume_path"); err == nil { - t.Error("volume_path ne devrait plus exister en DB") - } -} - -func TestStartVM_StorageReturnedInResponse(t *testing.T) { - s, db := newTestServer(t) - kv.AddInDB(db, "subnet/sn-1/state", "created") - kv.AddInDB(db, "subnet/sn-1/vpc", "vpc-1") - - body, _ := json.Marshal(VMCreateRequest{ - Name: "vm-11", - Interfaces: []VMInterface{ - {Subnet: "sn-1", IP: "10.0.0.11", Primary: true}, - }, - Storage: []VMStorage{ - {Path: "/data/root.qcow2", Dev: "sda"}, - }, - }) - - 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", w.Code) - } - - var vm VM - json.NewDecoder(w.Body).Decode(&vm) - if len(vm.Storage) != 1 { - t.Fatalf("attendu 1 disque dans la réponse, obtenu %d", len(vm.Storage)) - } - if vm.Storage[0].Dev != "sda" || vm.Storage[0].Path != "/data/root.qcow2" { - t.Errorf("disque inattendu dans la réponse : %+v", vm.Storage[0]) - } -} diff --git a/internal/api/agent/vms.go b/internal/api/agent/vms.go index 2ea5748..4dc6136 100644 --- a/internal/api/agent/vms.go +++ b/internal/api/agent/vms.go @@ -77,21 +77,16 @@ func (s *Server) startVM(w http.ResponseWriter, r *http.Request) { return } - disks := make([]dispatcher.VMDisk, len(req.Storage)) - for i, s := range req.Storage { - disks[i] = dispatcher.VMDisk{Path: s.Path, Dev: s.Dev} - } - cmd := dispatcher.StartVMCommand{ - Name: req.Name, - Subnet: primary.Subnet, - IP: primary.IP, - Disks: disks, - Memory: req.Memory, - CPUs: req.CPUs, - UEFI: req.UEFI, - Password: req.Password, - SSHKey: req.SSHKey, + Name: req.Name, + Subnet: primary.Subnet, + IP: primary.IP, + VolumePath: req.Storage[0].Path, + Memory: req.Memory, + CPUs: req.CPUs, + UEFI: req.UEFI, + Password: req.Password, + SSHKey: req.SSHKey, } if err := s.dispatcher.Prepare(cmd); err != nil { diff --git a/internal/dispatcher/agent/vm_commands.go b/internal/dispatcher/agent/vm_commands.go index bdfeaec..1c4bbce 100644 --- a/internal/dispatcher/agent/vm_commands.go +++ b/internal/dispatcher/agent/vm_commands.go @@ -13,21 +13,16 @@ import ( "github.com/dgraph-io/badger/v4" ) -type VMDisk struct { - Path string - Dev string -} - type StartVMCommand struct { - Name string - Subnet string - IP string - Disks []VMDisk - Memory int - CPUs int - UEFI bool - Password string - SSHKey string + Name string + Subnet string + IP string + VolumePath string + Memory int + CPUs int + UEFI bool + Password string + SSHKey string } func (c StartVMCommand) Prepare(db *badger.DB, _ *configuration.Config) error { @@ -49,9 +44,7 @@ func (c StartVMCommand) Prepare(db *badger.DB, _ *configuration.Config) error { kv.AddInDB(db, "vm/"+c.Name+"/subnet", c.Subnet) kv.AddInDB(db, "vm/"+c.Name+"/ip", c.IP) kv.AddInDB(db, "vm/"+c.Name+"/metadata_port", strconv.Itoa(port)) - for _, d := range c.Disks { - kv.AddInDB(db, "vm/"+c.Name+"/disk/"+d.Dev, d.Path) - } + kv.AddInDB(db, "vm/"+c.Name+"/volume_path", c.VolumePath) kv.AddInDB(db, "vm/"+c.Name+"/memory", strconv.Itoa(c.Memory)) kv.AddInDB(db, "vm/"+c.Name+"/cpus", strconv.Itoa(c.CPUs)) if c.UEFI { diff --git a/internal/dispatcher/agent/vm_commands_test.go b/internal/dispatcher/agent/vm_commands_test.go deleted file mode 100644 index ba400b9..0000000 --- a/internal/dispatcher/agent/vm_commands_test.go +++ /dev/null @@ -1,130 +0,0 @@ -package dispatcher - -import ( - "testing" - - "git.g3e.fr/syonad/two/pkg/db/kv" -) - -// --- StartVMCommand.Prepare : écriture des disques en DB --- - -func TestStartVMCommand_Prepare_SingleDisk(t *testing.T) { - _, db := newTestDispatcher(t) - kv.AddInDB(db, "subnet/sn-1/state", "created") - kv.AddInDB(db, "subnet/sn-1/vpc", "vpc-1") - - cmd := StartVMCommand{ - Name: "vm-1", - Subnet: "sn-1", - IP: "10.0.0.5", - Disks: []VMDisk{{Path: "/data/root.qcow2", Dev: "sda"}}, - } - if err := cmd.Prepare(db, nil); err != nil { - t.Fatalf("Prepare a échoué : %v", err) - } - - path, err := kv.GetFromDB(db, "vm/vm-1/disk/sda") - if err != nil { - t.Fatalf("clé disk/sda absente en DB : %v", err) - } - if path != "/data/root.qcow2" { - t.Errorf("path attendu /data/root.qcow2, obtenu %q", path) - } -} - -func TestStartVMCommand_Prepare_MultiDisk(t *testing.T) { - _, db := newTestDispatcher(t) - kv.AddInDB(db, "subnet/sn-1/state", "created") - kv.AddInDB(db, "subnet/sn-1/vpc", "vpc-1") - - cmd := StartVMCommand{ - Name: "vm-2", - Subnet: "sn-1", - IP: "10.0.0.6", - Disks: []VMDisk{ - {Path: "/data/root.qcow2", Dev: "sda"}, - {Path: "/data/data.qcow2", Dev: "sdb"}, - }, - } - if err := cmd.Prepare(db, nil); err != nil { - t.Fatalf("Prepare a échoué : %v", err) - } - - for dev, want := range map[string]string{ - "sda": "/data/root.qcow2", - "sdb": "/data/data.qcow2", - } { - got, err := kv.GetFromDB(db, "vm/vm-2/disk/"+dev) - if err != nil { - t.Fatalf("clé disk/%s absente en DB : %v", dev, err) - } - if got != want { - t.Errorf("disk/%s : attendu %q, obtenu %q", dev, want, got) - } - } -} - -func TestStartVMCommand_Prepare_SlotGap(t *testing.T) { - _, db := newTestDispatcher(t) - kv.AddInDB(db, "subnet/sn-1/state", "created") - kv.AddInDB(db, "subnet/sn-1/vpc", "vpc-1") - - // sdb absent au boot — slot réservé pour hotplug - cmd := StartVMCommand{ - Name: "vm-3", - Subnet: "sn-1", - IP: "10.0.0.7", - Disks: []VMDisk{ - {Path: "/data/root.qcow2", Dev: "sda"}, - {Path: "/data/extra.qcow2", Dev: "sdc"}, - }, - } - if err := cmd.Prepare(db, nil); err != nil { - t.Fatalf("Prepare a échoué : %v", err) - } - - if _, err := kv.GetFromDB(db, "vm/vm-3/disk/sda"); err != nil { - t.Fatalf("disk/sda absent : %v", err) - } - if _, err := kv.GetFromDB(db, "vm/vm-3/disk/sdc"); err != nil { - t.Fatalf("disk/sdc absent : %v", err) - } - if _, err := kv.GetFromDB(db, "vm/vm-3/disk/sdb"); err == nil { - t.Error("disk/sdb ne devrait pas exister en DB") - } -} - -func TestStartVMCommand_Prepare_NoVolumePath(t *testing.T) { - _, db := newTestDispatcher(t) - kv.AddInDB(db, "subnet/sn-1/state", "created") - kv.AddInDB(db, "subnet/sn-1/vpc", "vpc-1") - - cmd := StartVMCommand{ - Name: "vm-4", - Subnet: "sn-1", - IP: "10.0.0.8", - Disks: []VMDisk{{Path: "/data/root.qcow2", Dev: "sda"}}, - } - if err := cmd.Prepare(db, nil); err != nil { - t.Fatalf("Prepare a échoué : %v", err) - } - - if _, err := kv.GetFromDB(db, "vm/vm-4/volume_path"); err == nil { - t.Error("volume_path ne devrait plus être écrit en DB") - } -} - -func TestStartVMCommand_Prepare_Duplicate(t *testing.T) { - _, db := newTestDispatcher(t) - kv.AddInDB(db, "vm/vm-exist/state", "started") - - cmd := StartVMCommand{ - Name: "vm-exist", - Subnet: "sn-1", - IP: "10.0.0.9", - Disks: []VMDisk{{Path: "/data/root.qcow2", Dev: "sda"}}, - } - if err := cmd.Prepare(db, nil); err == nil { - t.Error("Prepare devrait échouer si la VM existe déjà") - } -} diff --git a/internal/qemu/config.go b/internal/qemu/config.go deleted file mode 100644 index 9b50add..0000000 --- a/internal/qemu/config.go +++ /dev/null @@ -1,20 +0,0 @@ -package qemu - -type DiskConfig struct { - Path string - Dev string -} - -type Config struct { - Name string - TapID int - Mac string - Disks []DiskConfig - Memory int - CPUs int - UEFICodePath string - UEFIVarsPath string - SerialDir string - MonitorDir string - QMPDir string -} diff --git a/internal/qemu/start_linux.go b/internal/qemu/start_linux.go index a7869a2..3f40e73 100644 --- a/internal/qemu/start_linux.go +++ b/internal/qemu/start_linux.go @@ -7,10 +7,22 @@ import ( "os" "os/exec" "path/filepath" - "sort" - "strings" ) +type Config struct { + Name string + TapID int + Mac string + VolumePath string + Memory int + CPUs int + UEFICodePath string + UEFIVarsPath string + SerialDir string + MonitorDir string + QMPDir string +} + func Start(cfg Config) error { memory := cfg.Memory if memory == 0 { @@ -52,47 +64,8 @@ func Start(cfg Config) error { ) } - hasScsi := false - for _, d := range cfg.Disks { - if strings.HasPrefix(d.Dev, "sd") { - hasScsi = true - break - } - } - if hasScsi { - args = append(args, "-device", "virtio-scsi-pci,id=scsi0") - } - - sorted := make([]DiskConfig, len(cfg.Disks)) - copy(sorted, cfg.Disks) - // vd* avant sd* : les disques virtio-blk bootent en premier. - // À lettre égale de type, ordre alphabétique. - sort.Slice(sorted, func(i, j int) bool { - iVirtio := strings.HasPrefix(sorted[i].Dev, "vd") - jVirtio := strings.HasPrefix(sorted[j].Dev, "vd") - if iVirtio != jVirtio { - return iVirtio - } - return sorted[i].Dev < sorted[j].Dev - }) - - for idx, d := range sorted { - bootindex := idx + 1 - if strings.HasPrefix(d.Dev, "sd") { - scsiID := int(d.Dev[2] - 'a') - args = append(args, - "-drive", fmt.Sprintf("file=%s,if=none,id=%s", d.Path, d.Dev), - "-device", fmt.Sprintf("scsi-hd,drive=%s,bus=scsi0.0,scsi-id=%d,bootindex=%d", d.Dev, scsiID, bootindex), - ) - } else { - args = append(args, - "-drive", fmt.Sprintf("file=%s,if=none,id=%s", d.Path, d.Dev), - "-device", fmt.Sprintf("virtio-blk-pci,drive=%s,bootindex=%d", d.Dev, bootindex), - ) - } - } - args = append(args, + "-drive", fmt.Sprintf("file=%s,if=virtio", cfg.VolumePath), "-netdev", fmt.Sprintf("tap,id=net0,ifname=tap%d,script=no,downscript=no", cfg.TapID), "-device", fmt.Sprintf("virtio-net-pci,netdev=net0,mac=%s", cfg.Mac), "-daemonize", diff --git a/internal/qemu/start_other.go b/internal/qemu/start_other.go index c28411c..c9c9192 100644 --- a/internal/qemu/start_other.go +++ b/internal/qemu/start_other.go @@ -2,9 +2,17 @@ package qemu -import ( - "errors" -) +import "errors" + +type Config struct { + Name, Mac, VolumePath string + TapID, Memory, CPUs int + UEFICodePath string + UEFIVarsPath string + SerialDir string + MonitorDir string + QMPDir string +} func Start(_ Config) error { return errors.New("vm: not supported on this platform") diff --git a/internal/vm/create.go b/internal/vm/create.go index 3e0d4a7..529c73c 100644 --- a/internal/vm/create.go +++ b/internal/vm/create.go @@ -52,16 +52,11 @@ func StartVM(db *badger.DB, name string, cfg *configuration.Config) error { return fmt.Errorf("start metadata: %w", err) } - qDisks := make([]qemu.DiskConfig, len(d.disks)) - for i, disk := range d.disks { - qDisks[i] = qemu.DiskConfig{Path: disk.path, Dev: disk.dev} - } - qcfg := qemu.Config{ Name: name, TapID: d.tapID, Mac: d.mac, - Disks: qDisks, + VolumePath: d.volumePath, Memory: d.memory, CPUs: d.cpus, SerialDir: cfg.QEMU.SerialDir, diff --git a/internal/vm/data.go b/internal/vm/data.go index 3435fb1..1e2644e 100644 --- a/internal/vm/data.go +++ b/internal/vm/data.go @@ -11,11 +11,6 @@ import ( "github.com/dgraph-io/badger/v4" ) -type diskEntry struct { - path string - dev string -} - type vmData struct { subnetName string vpcName string @@ -25,7 +20,7 @@ type vmData struct { ip string metadataPort string mac string - disks []diskEntry + volumePath string memory int cpus int uefi bool @@ -87,18 +82,11 @@ func loadVM(db *badger.DB, name string) (vmData, error) { } d.mac = mac - diskEntries, err := kv.ListByPrefix(db, "vm/"+name+"/disk/") + volumePath, err := kv.GetFromDB(db, "vm/"+name+"/volume_path") if err != nil { - return d, fmt.Errorf("list disks: %w", err) - } - if len(diskEntries) == 0 { - return d, fmt.Errorf("no disks found for vm %q", name) - } - diskPrefix := "vm/" + name + "/disk/" - for key, path := range diskEntries { - dev := strings.TrimPrefix(key, diskPrefix) - d.disks = append(d.disks, diskEntry{path: path, dev: dev}) + return d, fmt.Errorf("get volume_path: %w", err) } + d.volumePath = volumePath memoryStr, err := kv.GetFromDB(db, "vm/"+name+"/memory") if err != nil { diff --git a/internal/vm/delete.go b/internal/vm/delete.go index 3d808e6..7e47f3f 100644 --- a/internal/vm/delete.go +++ b/internal/vm/delete.go @@ -33,18 +33,25 @@ func StopVM(db *badger.DB, name string, cfg *configuration.Config) error { socketPath := filepath.Join(cfg.QEMU.QMPDir, name+".sock") - if _, err := os.Stat(socketPath); err == nil { - // socket présent : tenter l'arrêt gracieux - if _, err := qmp.Send(socketPath, []string{`{"execute":"system_powerdown"}`}); err == nil { - waitQMPDead(socketPath, - time.Duration(cfg.Dispatcher.TimeoutSeconds)*time.Second, - time.Duration(cfg.Dispatcher.PollSeconds)*time.Second, - ) - } - // connexion QMP échouée : QEMU déjà mort + if _, err := qmp.Send(socketPath, []string{`{"execute":"system_powerdown"}`}); err != nil { + return fmt.Errorf("qmp system_powerdown: %w", err) } - // socket absent ou QEMU déjà arrêté : cleanup direct + // attendre l'arrêt effectif de la VM ; forcer via quit après timeout + timeout := time.After(time.Duration(cfg.Dispatcher.TimeoutSeconds) * time.Second) + poll := time.Duration(cfg.Dispatcher.PollSeconds) * time.Second + stopped := false + for !stopped { + select { + case <-timeout: + qmp.Send(socketPath, []string{`{"execute":"quit"}`}) + stopped = true + case <-time.After(poll): + if _, err := qmp.Send(socketPath, nil); err != nil { + stopped = true + } + } + } if err := netns.Call(d.vpcName, func() error { return iptables.DeleteMetadataRedirect(d.ip, d.interfaceIP, d.metadataPort) @@ -67,18 +74,3 @@ func StopVM(db *badger.DB, name string, cfg *configuration.Config) error { return kv.AddInDB(db, "vm/"+name+"/state", "stopped") } - -func waitQMPDead(socketPath string, timeout, poll time.Duration) { - timer := time.After(timeout) - for { - select { - case <-timer: - qmp.Send(socketPath, []string{`{"execute":"quit"}`}) - return - case <-time.After(poll): - if _, err := qmp.Send(socketPath, nil); err != nil { - return - } - } - } -}