Compare commits

..

28 commits

Author SHA1 Message Date
a3b85f5926
Merge branch 'feature-32'
All checks were successful
Pre Release Workflow / set-release-target (push) Successful in 1s
Pre Release Workflow / build (metadata, amd64, linux) (push) Successful in 1m40s
Pre Release Workflow / build (agent, amd64, linux) (push) Successful in 2m1s
Pre Release Workflow / upload-scripts (run-dnsmasq-in-netns.sh) (push) Successful in 6s
Pre Release Workflow / prerelease (push) Successful in 10s
2026-06-14 23:05:16 +02:00
4ab880a32d
f-32: vm: fix incomplete stop vm
All checks were successful
Pre Release Workflow / set-release-target (push) Successful in 1s
Pre Release Workflow / build (agent, amd64, linux) (push) Successful in 1m36s
Pre Release Workflow / build (metadata, amd64, linux) (push) Successful in 2m19s
Pre Release Workflow / upload-scripts (run-dnsmasq-in-netns.sh) (push) Successful in 6s
Pre Release Workflow / prerelease (push) Successful in 11s
Signed-off-by: GnomeZworc <nicolas.boufidjeline@g3e.fr>
2026-05-24 19:02:31 +02:00
c3f26836cd
f-32: fix: update boot order
All checks were successful
Pre Release Workflow / set-release-target (push) Successful in 1s
Pre Release Workflow / upload-scripts (run-dnsmasq-in-netns.sh) (push) Successful in 6s
Pre Release Workflow / build (agent, amd64, linux) (push) Successful in 1m34s
Pre Release Workflow / build (metadata, amd64, linux) (push) Successful in 1m31s
Pre Release Workflow / prerelease (push) Successful in 10s
Signed-off-by: GnomeZworc <nicolas.boufidjeline@g3e.fr>
2026-05-24 16:47:49 +02:00
9111045417
f-32: test: add test for vms
All checks were successful
Pre Release Workflow / set-release-target (push) Successful in 1s
Pre Release Workflow / build (agent, amd64, linux) (push) Successful in 3m33s
Pre Release Workflow / build (metadata, amd64, linux) (push) Successful in 1m31s
Pre Release Workflow / upload-scripts (run-dnsmasq-in-netns.sh) (push) Successful in 6s
Pre Release Workflow / prerelease (push) Successful in 10s
Signed-off-by: GnomeZworc <nicolas.boufidjeline@g3e.fr>
2026-05-24 16:18:53 +02:00
a0637d827a
f-32: syntax: fix duplicated code
Signed-off-by: GnomeZworc <nicolas.boufidjeline@g3e.fr>
2026-05-24 16:13:11 +02:00
89899005d9
f-32: disk: add disk in qemu start
Signed-off-by: GnomeZworc <nicolas.boufidjeline@g3e.fr>
2026-05-24 15:56:45 +02:00
85c8c4e590
f-32: disk: add boot multidisk handle
Signed-off-by: GnomeZworc <nicolas.boufidjeline@g3e.fr>
2026-05-24 15:44:48 +02:00
437766d812
Merge branch 'feature-28' 2026-05-23 22:43:33 +02:00
2a5473eb22
f-28: qemu: execute with uefi vars
Signed-off-by: GnomeZworc <nicolas.boufidjeline@g3e.fr>
2026-05-23 22:30:32 +02:00
ec613996a0
f-28: vms: Add uefi param
Signed-off-by: GnomeZworc <nicolas.boufidjeline@g3e.fr>
2026-05-23 22:28:19 +02:00
325f1acff5
f-28: api: add uefi params
Signed-off-by: GnomeZworc <nicolas.boufidjeline@g3e.fr>
2026-05-23 22:26:21 +02:00
c602725de9
f-28: vm: add param in config file
Signed-off-by: GnomeZworc <nicolas.boufidjeline@g3e.fr>
2026-05-23 22:22:40 +02:00
732a293857
f-28: fix: debug somme minor errors
Signed-off-by: GnomeZworc <nicolas.boufidjeline@g3e.fr>
2026-05-21 00:15:45 +02:00
b41b4f2518
f-28: generate metadata_port automatically at vm creation
Signed-off-by: GnomeZworc <nicolas.boufidjeline@g3e.fr>
2026-05-18 23:28:18 +02:00
1b56a42627
f-28: refactor dhcp config to use VPCRoute and DefaultGateway
Signed-off-by: GnomeZworc <nicolas.boufidjeline@g3e.fr>
2026-05-18 23:06:52 +02:00
c0caf1a24c
f-28: add GetDefaultGateway to netif
Signed-off-by: GnomeZworc <nicolas.boufidjeline@g3e.fr>
2026-05-18 23:06:45 +02:00
bb5698fdda
f-28: add subnet default_route field
Signed-off-by: GnomeZworc <nicolas.boufidjeline@g3e.fr>
2026-05-18 23:06:41 +02:00
76a840b80a
f-28: add vpc cidr field
Signed-off-by: GnomeZworc <nicolas.boufidjeline@g3e.fr>
2026-05-18 23:06:36 +02:00
2504b435a3
f-28: fix: dhcp do not emit local default route
Signed-off-by: GnomeZworc <nicolas.boufidjeline@g3e.fr>
2026-05-18 23:06:29 +02:00
32b78a84f9
f-28: fix: add proper dhcp handle
Signed-off-by: GnomeZworc <nicolas.boufidjeline@g3e.fr>
2026-05-18 22:28:51 +02:00
848f965883
f-28: fix: renomage api param
Signed-off-by: GnomeZworc <nicolas.boufidjeline@g3e.fr>
2026-05-18 21:58:33 +02:00
9492de7a2b
f-28: test: add tests
Signed-off-by: GnomeZworc <nicolas.boufidjeline@g3e.fr>
2026-05-18 21:50:29 +02:00
91a5d7ac78
f-28: bridge: fix ebtables
Signed-off-by: GnomeZworc <nicolas.boufidjeline@g3e.fr>
2026-05-18 21:11:09 +02:00
d457c73198
f-28: bridge: first split bridge vxlan
Signed-off-by: GnomeZworc <nicolas.boufidjeline@g3e.fr>
2026-05-18 20:45:39 +02:00
9a16cf011a
f-28: api: implement new model usage
Signed-off-by: GnomeZworc <nicolas.boufidjeline@g3e.fr>
2026-05-18 20:15:16 +02:00
cef465da2e
f-28: api: update api comportement and model
Signed-off-by: GnomeZworc <nicolas.boufidjeline@g3e.fr>
2026-05-18 20:11:33 +02:00
7065a9f431
Merge branch 'feature-25' 2026-05-18 19:44:33 +02:00
0506be8a87
f-25: code: add db api
Signed-off-by: GnomeZworc <nicolas.boufidjeline@g3e.fr>
2026-04-30 23:29:07 +02:00
38 changed files with 1321 additions and 267 deletions

View file

@ -294,13 +294,17 @@ components:
VPCCreateRequest:
type: object
required: [name]
required: [name, cidr]
properties:
name:
type: string
description: Unique name for the VPC, must follow the format vp-[id]
pattern: '^vp-.+'
example: vp-00001
cidr:
type: string
description: CIDR block for the entire VPC address space
example: "10.0.0.0/16"
VPC:
type: object
@ -312,10 +316,13 @@ components:
type: string
enum: [creating, created, deleting, deleted]
example: created
cidr:
type: string
example: "10.0.0.0/16"
SubnetCreateRequest:
type: object
required: [name, vpc, vxlan_id, gateway_ip, cidr]
required: [name, vpc, interface_ip, cidr]
properties:
name:
type: string
@ -325,15 +332,24 @@ components:
type: string
description: Parent VPC name
example: vpc1
mode:
type: string
description: >
Subnet mode. "vxlan" (default): creates a VXLAN tunnel and a host bridge.
"bridge": attaches directly to an existing bridge resolved from iface_type in the agent config.
"vlan" is reserved for future use.
enum: [vxlan, bridge]
default: vxlan
example: vxlan
vxlan_id:
type: integer
description: VXLAN VNI identifier
description: VXLAN VNI identifier. Required when mode is "vxlan", ignored otherwise.
example: 100
iface_type:
type: string
description: Interface type key defined in the agent config (e.g. vms, internet, admin). Falls back to default_interface if omitted or unknown.
example: vms
gateway_ip:
interface_ip:
type: string
format: ipv4
description: Gateway IP for the subnet
@ -342,6 +358,12 @@ components:
type: string
description: Subnet CIDR block
example: "10.10.10.0/24"
default_route:
type: boolean
description: >
If true, advertise a default route via DHCP. For vxlan mode the gateway is the interface IP.
For bridge mode the gateway is read from the host routing table.
default: false
Subnet:
type: object
@ -356,30 +378,35 @@ components:
vpc:
type: string
example: vpc1
mode:
type: string
enum: [vxlan, bridge]
example: vxlan
vxlan_id:
type: integer
description: VXLAN VNI. Present only when mode is "vxlan".
example: 100
local_iface:
type: string
description: Resolved interface name
description: Resolved interface name from agent config
example: br-000000
gateway_ip:
interface_ip:
type: string
example: "10.10.10.1"
cidr:
type: string
example: "10.10.10.0/24"
default_route:
type: boolean
example: false
VMCreateRequest:
type: object
required: [name, metadata_port, interfaces, storage]
required: [name, interfaces, storage]
properties:
name:
type: string
example: vm-00001
metadata_port:
type: string
example: "80"
memory:
type: integer
description: Memory in MB (default 512)
@ -403,6 +430,10 @@ components:
minItems: 1
items:
$ref: "#/components/schemas/VMStorage"
uefi:
type: boolean
description: Boot with UEFI firmware (OVMF). Defaults to false (SeaBIOS).
example: false
VMInterface:
type: object
@ -460,6 +491,9 @@ components:
type: array
items:
$ref: "#/components/schemas/VMStorage"
uefi:
type: boolean
example: false
Error:
type: object

View file

@ -51,6 +51,10 @@ func main() {
d := dispatcher.New(q, db, cfg, log.With(slog.String("component", "dispatcher")))
go agentapi.New(d, db, log.With(slog.String("component", "api"))).Start(apiAddr)
go promserver.Start(promAddr, registry)
if cfg.Admin.Enabled {
adminAddr := fmt.Sprintf("%s:%d", cfg.Admin.Address, cfg.Admin.Port)
go kv.NewAdminServer(db, log.With(slog.String("component", "admin"))).Start(adminAddr)
}
select {}
}

View file

@ -39,6 +39,24 @@ interfaces:
metadata:
run_dir: "/run/two/metadata"
# QEMU runtime paths
qemu:
# UEFI firmware (requires apt install ovmf on Debian/Ubuntu)
ovmf_code_path: "/usr/share/OVMF/OVMF_CODE.fd"
ovmf_vars_template: "/usr/share/OVMF/OVMF_VARS.fd"
# Per-VM UEFI variable store (writable copy, created at start / deleted at stop)
uefi_vars_dir: "/run/two/vms/efi"
# QEMU Unix socket directories
serial_dir: "/run/two/vms/serial"
monitor_dir: "/run/two/vms/monitor"
qmp_dir: "/run/two/vms/qmp"
# Admin API (read-only DB inspection, loopback only)
admin:
enabled: false
address: "127.0.0.1"
port: 9091
# Logging configuration
logger:
# Log level: debug, info, warn, error (default: info)

View file

@ -2,30 +2,36 @@ package agentapi
type VPCCreateRequest struct {
Name string `json:"name"`
CIDR string `json:"cidr"`
}
type VPC struct {
Name string `json:"name"`
State string `json:"state"`
CIDR string `json:"cidr"`
}
type SubnetCreateRequest struct {
Name string `json:"name"`
VPC string `json:"vpc"`
VxlanID int `json:"vxlan_id"`
IfaceType string `json:"iface_type"`
GatewayIP string `json:"gateway_ip"`
CIDR string `json:"cidr"`
Name string `json:"name"`
VPC string `json:"vpc"`
Mode string `json:"mode"`
VxlanID int `json:"vxlan_id"`
IfaceType string `json:"iface_type"`
InterfaceIP string `json:"interface_ip"`
CIDR string `json:"cidr"`
DefaultRoute bool `json:"default_route"`
}
type Subnet struct {
Name string `json:"name"`
State string `json:"state"`
VPC string `json:"vpc"`
VxlanID int `json:"vxlan_id"`
LocalIface string `json:"local_iface"`
GatewayIP string `json:"gateway_ip"`
CIDR string `json:"cidr"`
Name string `json:"name"`
State string `json:"state"`
VPC string `json:"vpc"`
Mode string `json:"mode"`
VxlanID int `json:"vxlan_id"`
LocalIface string `json:"local_iface"`
InterfaceIP string `json:"interface_ip"`
CIDR string `json:"cidr"`
DefaultRoute bool `json:"default_route"`
}
type VMInterface struct {
@ -40,14 +46,14 @@ type VMStorage struct {
}
type VMCreateRequest struct {
Name string `json:"name"`
MetadataPort string `json:"metadata_port"`
Memory int `json:"memory"`
CPUs int `json:"cpus"`
Password string `json:"password"`
SSHKey string `json:"sshkey"`
Interfaces []VMInterface `json:"interfaces"`
Storage []VMStorage `json:"storage"`
Name string `json:"name"`
Memory int `json:"memory"`
CPUs int `json:"cpus"`
UEFI bool `json:"uefi"`
Password string `json:"password"`
SSHKey string `json:"sshkey"`
Interfaces []VMInterface `json:"interfaces"`
Storage []VMStorage `json:"storage"`
}
type VM struct {
@ -56,6 +62,7 @@ type VM struct {
MetadataPort string `json:"metadata_port"`
Memory int `json:"memory"`
CPUs int `json:"cpus"`
UEFI bool `json:"uefi"`
Interfaces []VMInterface `json:"interfaces"`
Storage []VMStorage `json:"storage"`
}

View file

@ -47,14 +47,18 @@ func (s *Server) getSubnet(w http.ResponseWriter, _ *http.Request, name string)
sub.State = value
case "vpc":
sub.VPC = value
case "mode":
sub.Mode = value
case "vxlan_id":
sub.VxlanID, _ = strconv.Atoi(value)
case "local_iface":
sub.LocalIface = value
case "gateway_ip":
sub.GatewayIP = value
case "interface_ip":
sub.InterfaceIP = value
case "cidr":
sub.CIDR = value
case "default_route":
sub.DefaultRoute = value == "true"
}
}
w.WriteHeader(http.StatusOK)

View file

@ -60,7 +60,7 @@ func TestPostSubnet_Created(t *testing.T) {
Name: "sn-new",
VPC: "vpc-1",
IfaceType: "vms",
GatewayIP: "10.0.0.1",
InterfaceIP: "10.0.0.1",
CIDR: "10.0.0.0/24",
}
body, _ := json.Marshal(req)
@ -95,7 +95,7 @@ func TestPostSubnet_IfaceTypeOptional(t *testing.T) {
req := SubnetCreateRequest{
Name: "sn-opt",
VPC: "vpc-1",
GatewayIP: "10.0.0.1",
InterfaceIP: "10.0.0.1",
CIDR: "10.0.0.0/24",
// IfaceType omis — doit utiliser default_interface
}
@ -113,7 +113,7 @@ func TestPostSubnet_VPCNotFound(t *testing.T) {
Name: "sn-1",
VPC: "vpc-inexistant",
IfaceType: "vms",
GatewayIP: "10.0.0.1",
InterfaceIP: "10.0.0.1",
CIDR: "10.0.0.0/24",
}
body, _ := json.Marshal(req)
@ -132,7 +132,7 @@ func TestPostSubnet_Duplicate(t *testing.T) {
Name: "sn-exist",
VPC: "vpc-1",
IfaceType: "vms",
GatewayIP: "10.0.0.1",
InterfaceIP: "10.0.0.1",
CIDR: "10.0.0.0/24",
}
body, _ := json.Marshal(req)
@ -150,7 +150,52 @@ func TestPostSubnet_VPCDeleting(t *testing.T) {
Name: "sn-1",
VPC: "vpc-dying",
IfaceType: "vms",
GatewayIP: "10.0.0.1",
InterfaceIP: "10.0.0.1",
CIDR: "10.0.0.0/24",
}
body, _ := json.Marshal(req)
w := httptest.NewRecorder()
s.SubnetsHandler(w, httptest.NewRequest(http.MethodPost, "/subnets", bytes.NewReader(body)))
if w.Code != http.StatusUnprocessableEntity {
t.Errorf("attendu 422, obtenu %d", w.Code)
}
}
func TestPostSubnet_BridgeMode_Success(t *testing.T) {
s, db := newTestServer(t)
kv.AddInDB(db, "vpc/vpc-1/state", "created")
req := SubnetCreateRequest{
Name: "sn-br",
VPC: "vpc-1",
Mode: "bridge",
IfaceType: "vms",
InterfaceIP: "10.0.0.1",
CIDR: "10.0.0.0/24",
}
body, _ := json.Marshal(req)
w := httptest.NewRecorder()
s.SubnetsHandler(w, httptest.NewRequest(http.MethodPost, "/subnets", bytes.NewReader(body)))
if w.Code != http.StatusAccepted {
t.Fatalf("attendu 202, obtenu %d: %s", w.Code, w.Body.String())
}
var result Subnet
json.NewDecoder(w.Body).Decode(&result)
if result.Mode != "bridge" {
t.Errorf("mode attendu bridge, obtenu %q", result.Mode)
}
if result.VxlanID != 0 {
t.Errorf("vxlan_id devrait être 0 en mode bridge, obtenu %d", result.VxlanID)
}
}
func TestPostSubnet_UnknownMode(t *testing.T) {
s, db := newTestServer(t)
kv.AddInDB(db, "vpc/vpc-1/state", "created")
req := SubnetCreateRequest{
Name: "sn-1",
VPC: "vpc-1",
Mode: "vlan",
InterfaceIP: "10.0.0.1",
CIDR: "10.0.0.0/24",
}
body, _ := json.Marshal(req)
@ -177,7 +222,7 @@ func TestGetSubnet_Found(t *testing.T) {
kv.AddInDB(db, "subnet/sn-1/state", "created")
kv.AddInDB(db, "subnet/sn-1/vpc", "vpc-1")
kv.AddInDB(db, "subnet/sn-1/cidr", "10.0.0.0/24")
kv.AddInDB(db, "subnet/sn-1/gateway_ip", "10.0.0.1")
kv.AddInDB(db, "subnet/sn-1/interface_ip", "10.0.0.1")
req := httptest.NewRequest(http.MethodGet, "/subnets/sn-1", nil)
w := httptest.NewRecorder()
s.SubnetByNameHandler(w, req)

View file

@ -44,14 +44,18 @@ func (s *Server) listSubnets(w http.ResponseWriter, _ *http.Request) {
subnets[name].State = value
case "vpc":
subnets[name].VPC = value
case "mode":
subnets[name].Mode = value
case "vxlan_id":
subnets[name].VxlanID, _ = strconv.Atoi(value)
case "local_iface":
subnets[name].LocalIface = value
case "gateway_ip":
subnets[name].GatewayIP = value
case "interface_ip":
subnets[name].InterfaceIP = value
case "cidr":
subnets[name].CIDR = value
case "default_route":
subnets[name].DefaultRoute = value == "true"
}
}
result := make([]Subnet, 0, len(subnets))
@ -69,18 +73,20 @@ func (s *Server) postSubnet(w http.ResponseWriter, r *http.Request) {
json.NewEncoder(w).Encode(ErrorResponse{Error: "invalid request body"})
return
}
if req.Name == "" || req.VPC == "" || req.GatewayIP == "" || req.CIDR == "" {
if req.Name == "" || req.VPC == "" || req.InterfaceIP == "" || req.CIDR == "" {
w.WriteHeader(http.StatusBadRequest)
json.NewEncoder(w).Encode(ErrorResponse{Error: "name, vpc, gateway_ip and cidr are required"})
json.NewEncoder(w).Encode(ErrorResponse{Error: "name, vpc, interface_ip and cidr are required"})
return
}
cmd := dispatcher.CreateSubnetCommand{
Name: req.Name,
VPC: req.VPC,
VxlanID: req.VxlanID,
IfaceType: req.IfaceType,
GatewayIP: req.GatewayIP,
CIDR: req.CIDR,
Name: req.Name,
VPC: req.VPC,
Mode: req.Mode,
VxlanID: req.VxlanID,
IfaceType: req.IfaceType,
InterfaceIP: req.InterfaceIP,
CIDR: req.CIDR,
DefaultRoute: req.DefaultRoute,
}
if err := s.dispatcher.Prepare(cmd); err != nil {
if _, dbErr := kv.GetFromDB(s.db, "subnet/"+req.Name+"/state"); dbErr == nil {
@ -109,14 +115,18 @@ func (s *Server) postSubnet(w http.ResponseWriter, r *http.Request) {
sub.State = value
case "vpc":
sub.VPC = value
case "mode":
sub.Mode = value
case "vxlan_id":
sub.VxlanID, _ = strconv.Atoi(value)
case "local_iface":
sub.LocalIface = value
case "gateway_ip":
sub.GatewayIP = value
case "interface_ip":
sub.InterfaceIP = value
case "cidr":
sub.CIDR = value
case "default_route":
sub.DefaultRoute = value == "true"
}
}
w.WriteHeader(http.StatusAccepted)

View file

@ -79,6 +79,7 @@ func vmFromDB(name string, entries map[string]string) (VM, error) {
vm.MetadataPort = entries[prefix+"metadata_port"]
vm.Memory, _ = strconv.Atoi(entries[prefix+"memory"])
vm.CPUs, _ = strconv.Atoi(entries[prefix+"cpus"])
vm.UEFI = entries[prefix+"uefi"] == "true"
subnet := entries[prefix+"subnet"]
ip := entries[prefix+"ip"]
@ -86,8 +87,11 @@ func vmFromDB(name string, entries map[string]string) (VM, error) {
vm.Interfaces = []VMInterface{{Subnet: subnet, IP: ip, Primary: true}}
}
if path := entries[prefix+"volume_path"]; path != "" {
vm.Storage = []VMStorage{{Path: path}}
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})
}
}
return vm, nil

View file

@ -0,0 +1,158 @@
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])
}
}

View file

@ -58,9 +58,9 @@ func (s *Server) startVM(w http.ResponseWriter, r *http.Request) {
json.NewEncoder(w).Encode(ErrorResponse{Error: "invalid request body"})
return
}
if req.Name == "" || req.MetadataPort == "" || len(req.Interfaces) == 0 || len(req.Storage) == 0 {
if req.Name == "" || len(req.Interfaces) == 0 || len(req.Storage) == 0 {
w.WriteHeader(http.StatusBadRequest)
json.NewEncoder(w).Encode(ErrorResponse{Error: "name, metadata_port, interfaces and storage are required"})
json.NewEncoder(w).Encode(ErrorResponse{Error: "name, interfaces and storage are required"})
return
}
@ -77,16 +77,21 @@ 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,
MetadataPort: req.MetadataPort,
VolumePath: req.Storage[0].Path,
Memory: req.Memory,
CPUs: req.CPUs,
Password: req.Password,
SSHKey: req.SSHKey,
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,
}
if err := s.dispatcher.Prepare(cmd); err != nil {

View file

@ -35,8 +35,9 @@ func (s *Server) getVpc(w http.ResponseWriter, _ *http.Request, name string) {
json.NewEncoder(w).Encode(ErrorResponse{Error: "vpc not found"})
return
}
cidr, _ := kv.GetFromDB(s.db, "vpc/"+name+"/cidr")
w.WriteHeader(http.StatusOK)
json.NewEncoder(w).Encode(VPC{Name: name, State: state})
json.NewEncoder(w).Encode(VPC{Name: name, State: state, CIDR: cidr})
}
func (s *Server) deleteVpc(w http.ResponseWriter, _ *http.Request, name string) {

View file

@ -53,7 +53,7 @@ func TestListVpcs_InvalidMethod(t *testing.T) {
func TestPostVpc_Created(t *testing.T) {
s, _ := newTestServer(t)
body, _ := json.Marshal(VPCCreateRequest{Name: "vpc-new"})
body, _ := json.Marshal(VPCCreateRequest{Name: "vpc-new", CIDR: "10.0.0.0/16"})
w := httptest.NewRecorder()
s.VpcsHandler(w, httptest.NewRequest(http.MethodPost, "/vpcs", bytes.NewReader(body)))
if w.Code != http.StatusAccepted {
@ -67,11 +67,34 @@ func TestPostVpc_Created(t *testing.T) {
if result.State != "creating" {
t.Errorf("state attendu creating, obtenu %q", result.State)
}
if result.CIDR != "10.0.0.0/16" {
t.Errorf("cidr attendu 10.0.0.0/16, obtenu %q", result.CIDR)
}
}
func TestPostVpc_MissingName(t *testing.T) {
s, _ := newTestServer(t)
body, _ := json.Marshal(VPCCreateRequest{})
body, _ := json.Marshal(VPCCreateRequest{CIDR: "10.0.0.0/16"})
w := httptest.NewRecorder()
s.VpcsHandler(w, httptest.NewRequest(http.MethodPost, "/vpcs", bytes.NewReader(body)))
if w.Code != http.StatusBadRequest {
t.Errorf("attendu 400, obtenu %d", w.Code)
}
}
func TestPostVpc_MissingCIDR(t *testing.T) {
s, _ := newTestServer(t)
body, _ := json.Marshal(VPCCreateRequest{Name: "vpc-new"})
w := httptest.NewRecorder()
s.VpcsHandler(w, httptest.NewRequest(http.MethodPost, "/vpcs", bytes.NewReader(body)))
if w.Code != http.StatusBadRequest {
t.Errorf("attendu 400, obtenu %d", w.Code)
}
}
func TestPostVpc_InvalidCIDR(t *testing.T) {
s, _ := newTestServer(t)
body, _ := json.Marshal(VPCCreateRequest{Name: "vpc-new", CIDR: "not-a-cidr"})
w := httptest.NewRecorder()
s.VpcsHandler(w, httptest.NewRequest(http.MethodPost, "/vpcs", bytes.NewReader(body)))
if w.Code != http.StatusBadRequest {
@ -82,7 +105,7 @@ func TestPostVpc_MissingName(t *testing.T) {
func TestPostVpc_Duplicate(t *testing.T) {
s, db := newTestServer(t)
kv.AddInDB(db, "vpc/vpc-exist/state", "created")
body, _ := json.Marshal(VPCCreateRequest{Name: "vpc-exist"})
body, _ := json.Marshal(VPCCreateRequest{Name: "vpc-exist", CIDR: "10.0.0.0/16"})
w := httptest.NewRecorder()
s.VpcsHandler(w, httptest.NewRequest(http.MethodPost, "/vpcs", bytes.NewReader(body)))
if w.Code != http.StatusConflict {

View file

@ -2,6 +2,7 @@ package agentapi
import (
"encoding/json"
"net"
"net/http"
"strings"
@ -38,8 +39,11 @@ func (s *Server) listVpcs(w http.ResponseWriter, _ *http.Request) {
if _, ok := vpcs[name]; !ok {
vpcs[name] = &VPC{Name: name}
}
if parts[2] == "state" {
switch parts[2] {
case "state":
vpcs[name].State = value
case "cidr":
vpcs[name].CIDR = value
}
}
result := make([]VPC, 0, len(vpcs))
@ -62,7 +66,17 @@ func (s *Server) postVpc(w http.ResponseWriter, r *http.Request) {
json.NewEncoder(w).Encode(ErrorResponse{Error: "name is required"})
return
}
cmd := dispatcher.CreateVPCCommand{Name: req.Name}
if req.CIDR == "" {
w.WriteHeader(http.StatusBadRequest)
json.NewEncoder(w).Encode(ErrorResponse{Error: "cidr is required"})
return
}
if _, _, err := net.ParseCIDR(req.CIDR); err != nil {
w.WriteHeader(http.StatusBadRequest)
json.NewEncoder(w).Encode(ErrorResponse{Error: "invalid cidr"})
return
}
cmd := dispatcher.CreateVPCCommand{Name: req.Name, CIDR: req.CIDR}
if err := s.dispatcher.Prepare(cmd); err != nil {
w.WriteHeader(http.StatusConflict)
json.NewEncoder(w).Encode(ErrorResponse{Error: err.Error()})
@ -76,5 +90,5 @@ func (s *Server) postVpc(w http.ResponseWriter, r *http.Request) {
return
}
w.WriteHeader(http.StatusAccepted)
json.NewEncoder(w).Encode(VPC{Name: req.Name, State: state})
json.NewEncoder(w).Encode(VPC{Name: req.Name, State: state, CIDR: req.CIDR})
}

View file

@ -31,6 +31,19 @@ type Config struct {
Metadata struct {
RunDir string `mapstructure:"run_dir"`
} `mapstructure:"metadata"`
Admin struct {
Enabled bool `mapstructure:"enabled"`
Address string `mapstructure:"address"`
Port int `mapstructure:"port"`
} `mapstructure:"admin"`
QEMU struct {
OVMFCodePath string `mapstructure:"ovmf_code_path"`
OVMFVarsTemplate string `mapstructure:"ovmf_vars_template"`
UEFIVarsDir string `mapstructure:"uefi_vars_dir"`
SerialDir string `mapstructure:"serial_dir"`
MonitorDir string `mapstructure:"monitor_dir"`
QMPDir string `mapstructure:"qmp_dir"`
} `mapstructure:"qemu"`
DefaultInterface string `mapstructure:"default_interface"`
Interfaces map[string]string `mapstructure:"interfaces"`
}
@ -50,6 +63,15 @@ func LoadConfig(path string) (*Config, error) {
v.SetDefault("dispatcher.timeout_seconds", 300)
v.SetDefault("dispatcher.poll_seconds", 2)
v.SetDefault("metadata.run_dir", "/run/two/metadata")
v.SetDefault("qemu.ovmf_code_path", "/usr/share/OVMF/OVMF_CODE.fd")
v.SetDefault("qemu.ovmf_vars_template", "/usr/share/OVMF/OVMF_VARS.fd")
v.SetDefault("qemu.uefi_vars_dir", "/run/two/vms/uefi")
v.SetDefault("qemu.serial_dir", "/run/two/vms/serial")
v.SetDefault("qemu.monitor_dir", "/run/two/vms/monitor")
v.SetDefault("qemu.qmp_dir", "/run/two/vms/qmp")
v.SetDefault("admin.enabled", false)
v.SetDefault("admin.address", "127.0.0.1")
v.SetDefault("admin.port", 9091)
v.SetDefault("default_interface", "br-000000")
v.SetDefault("logger.level", "info")
v.SetDefault("logger.debug", false)

View file

@ -51,11 +51,15 @@ func TestIncrementIP_Carry(t *testing.T) {
func newConf(t *testing.T, cidr string) Config {
t.Helper()
_, network, _ := net.ParseCIDR(cidr)
_, vpcNet, _ := net.ParseCIDR("10.0.0.0/16")
gw := net.ParseIP("192.168.1.1").To4()
return Config{
Network: network,
Gateway: net.ParseIP("192.168.1.1").To4(),
Name: "test",
ConfDir: t.TempDir(),
Network: network,
VPCGateway: gw,
VPCRoute: vpcNet,
DefaultGateway: gw,
Name: "test",
ConfDir: t.TempDir(),
}
}
@ -84,13 +88,45 @@ func TestGenerateConfig_FilenameMatchesName(t *testing.T) {
}
}
func TestGenerateConfig_ContainsGateway(t *testing.T) {
func TestGenerateConfig_ContainsDefaultGateway(t *testing.T) {
conf := newConf(t, "192.168.1.0/29")
path, _, _ := GenerateConfig(conf)
content, _ := os.ReadFile(path)
if !strings.Contains(string(content), "dhcp-option=3,192.168.1.1") {
t.Errorf("gateway absente du fichier généré :\n%s", content)
t.Errorf("dhcp-option=3 absente du fichier généré :\n%s", content)
}
}
func TestGenerateConfig_NoDefaultGateway(t *testing.T) {
conf := newConf(t, "192.168.1.0/29")
conf.DefaultGateway = nil
path, _, _ := GenerateConfig(conf)
content, _ := os.ReadFile(path)
if strings.Contains(string(content), "dhcp-option=3,") {
t.Errorf("dhcp-option=3 présente alors que DefaultGateway=nil :\n%s", content)
}
}
func TestGenerateConfig_ContainsVPCRoute(t *testing.T) {
conf := newConf(t, "192.168.1.0/29")
path, _, _ := GenerateConfig(conf)
content, _ := os.ReadFile(path)
if !strings.Contains(string(content), "dhcp-option=121,10.0.0.0/16,192.168.1.1") {
t.Errorf("dhcp-option=121 absente ou incorrecte :\n%s", content)
}
}
func TestGenerateConfig_NoVPCRoute(t *testing.T) {
conf := newConf(t, "192.168.1.0/29")
conf.VPCRoute = nil
path, _, _ := GenerateConfig(conf)
content, _ := os.ReadFile(path)
if strings.Contains(string(content), "dhcp-option=121,") {
t.Errorf("dhcp-option=121 présente alors que VPCRoute=nil :\n%s", content)
}
}
@ -98,7 +134,6 @@ func TestGenerateConfig_ContainsDhcpRange(t *testing.T) {
_, network, _ := net.ParseCIDR("10.10.0.0/24")
conf := Config{
Network: network,
Gateway: net.ParseIP("10.10.0.1").To4(),
Name: "vpc1",
ConfDir: t.TempDir(),
}
@ -144,7 +179,6 @@ func TestGenerateConfig_CreatesConfDir(t *testing.T) {
_, network, _ := net.ParseCIDR("10.0.0.0/30")
conf := Config{
Network: network,
Gateway: net.ParseIP("10.0.0.1").To4(),
Name: "net",
ConfDir: dir,
}

View file

@ -14,7 +14,12 @@ func GenerateConfig(c Config) (string, map[string]string, error) {
var sb strings.Builder
fmt.Fprintf(&sb, "no-resolv\n")
fmt.Fprintf(&sb, "dhcp-range=%s,static,%s,12h\n", c.Network.IP.String(), mask)
fmt.Fprintf(&sb, "dhcp-option=3,%s\n", c.Gateway.String())
if c.VPCRoute != nil {
fmt.Fprintf(&sb, "dhcp-option=121,%s,%s\n", c.VPCRoute.String(), c.VPCGateway.String())
}
if c.DefaultGateway != nil {
fmt.Fprintf(&sb, "dhcp-option=3,%s\n", c.DefaultGateway.String())
}
fmt.Fprintf(&sb, "dhcp-option=6,1.1.1.1,8.8.8.8\n\n")
entries := make(map[string]string)

View file

@ -5,8 +5,10 @@ import (
)
type Config struct {
Network *net.IPNet
Gateway net.IP
Name string
ConfDir string
Network *net.IPNet
VPCGateway net.IP // next-hop for VPCRoute (option 121)
VPCRoute *net.IPNet // if non-nil, emit dhcp-option=121,VPCRoute,VPCGateway
DefaultGateway net.IP // if non-nil, emit dhcp-option=3,DefaultGateway
Name string
ConfDir string
}

View file

@ -12,15 +12,23 @@ import (
)
type CreateSubnetCommand struct {
Name string
VPC string
VxlanID int
IfaceType string
GatewayIP string
CIDR string
Name string
VPC string
Mode string
VxlanID int
IfaceType string
InterfaceIP string
CIDR string
DefaultRoute bool
}
func (c CreateSubnetCommand) Prepare(db *badger.DB, cfg *configuration.Config) error {
if c.Mode == "" {
c.Mode = "vxlan"
}
if c.Mode != "vxlan" && c.Mode != "bridge" {
return fmt.Errorf("unknown subnet mode %q", c.Mode)
}
if _, err := kv.GetFromDB(db, "subnet/"+c.Name+"/state"); err == nil {
return fmt.Errorf("subnet %q already exists", c.Name)
}
@ -37,10 +45,14 @@ func (c CreateSubnetCommand) Prepare(db *badger.DB, cfg *configuration.Config) e
}
kv.AddInDB(db, "subnet/"+c.Name+"/state", "creating")
kv.AddInDB(db, "subnet/"+c.Name+"/vpc", c.VPC)
kv.AddInDB(db, "subnet/"+c.Name+"/vxlan_id", strconv.Itoa(c.VxlanID))
kv.AddInDB(db, "subnet/"+c.Name+"/mode", c.Mode)
kv.AddInDB(db, "subnet/"+c.Name+"/local_iface", localIface)
kv.AddInDB(db, "subnet/"+c.Name+"/gateway_ip", c.GatewayIP)
kv.AddInDB(db, "subnet/"+c.Name+"/interface_ip", c.InterfaceIP)
kv.AddInDB(db, "subnet/"+c.Name+"/cidr", c.CIDR)
kv.AddInDB(db, "subnet/"+c.Name+"/default_route", strconv.FormatBool(c.DefaultRoute))
if c.Mode == "vxlan" {
kv.AddInDB(db, "subnet/"+c.Name+"/vxlan_id", strconv.Itoa(c.VxlanID))
}
return nil
}

View file

@ -20,7 +20,7 @@ func TestCreateSubnetCommand_Prepare_Success(t *testing.T) {
kv.AddInDB(db, "vpc/vpc-1/state", "created")
cmd := CreateSubnetCommand{
Name: "sn-1", VPC: "vpc-1", VxlanID: 100,
IfaceType: "vms", GatewayIP: "10.0.0.1", CIDR: "10.0.0.0/24",
IfaceType: "vms", InterfaceIP: "10.0.0.1", CIDR: "10.0.0.0/24",
}
if err := cmd.Prepare(db, testCfg()); err != nil {
t.Fatalf("Prepare a échoué : %v", err)
@ -40,7 +40,7 @@ func TestCreateSubnetCommand_Prepare_UsesIfaceTypeMapping(t *testing.T) {
kv.AddInDB(db, "vpc/vpc-1/state", "created")
cmd := CreateSubnetCommand{
Name: "sn-1", VPC: "vpc-1", VxlanID: 100,
IfaceType: "vms", GatewayIP: "10.0.0.1", CIDR: "10.0.0.0/24",
IfaceType: "vms", InterfaceIP: "10.0.0.1", CIDR: "10.0.0.0/24",
}
cmd.Prepare(db, testCfg())
iface, _ := kv.GetFromDB(db, "subnet/sn-1/local_iface")
@ -54,7 +54,7 @@ func TestCreateSubnetCommand_Prepare_UsesDefaultIfaceWhenTypeUnknown(t *testing.
kv.AddInDB(db, "vpc/vpc-1/state", "created")
cmd := CreateSubnetCommand{
Name: "sn-1", VPC: "vpc-1", VxlanID: 100,
IfaceType: "inconnu", GatewayIP: "10.0.0.1", CIDR: "10.0.0.0/24",
IfaceType: "inconnu", InterfaceIP: "10.0.0.1", CIDR: "10.0.0.0/24",
}
cmd.Prepare(db, testCfg())
iface, _ := kv.GetFromDB(db, "subnet/sn-1/local_iface")
@ -69,7 +69,7 @@ func TestCreateSubnetCommand_Prepare_Duplicate(t *testing.T) {
kv.AddInDB(db, "subnet/sn-exist/state", "created")
cmd := CreateSubnetCommand{
Name: "sn-exist", VPC: "vpc-1", VxlanID: 100,
IfaceType: "vms", GatewayIP: "10.0.0.1", CIDR: "10.0.0.0/24",
IfaceType: "vms", InterfaceIP: "10.0.0.1", CIDR: "10.0.0.0/24",
}
if err := cmd.Prepare(db, testCfg()); err == nil {
t.Error("Prepare devrait échouer sur un subnet déjà existant")
@ -80,7 +80,7 @@ func TestCreateSubnetCommand_Prepare_VPCNotFound(t *testing.T) {
_, db := newTestDispatcher(t)
cmd := CreateSubnetCommand{
Name: "sn-1", VPC: "vpc-inexistant", VxlanID: 100,
IfaceType: "vms", GatewayIP: "10.0.0.1", CIDR: "10.0.0.0/24",
IfaceType: "vms", InterfaceIP: "10.0.0.1", CIDR: "10.0.0.0/24",
}
if err := cmd.Prepare(db, testCfg()); err == nil {
t.Error("Prepare devrait échouer si le VPC n'existe pas")
@ -92,7 +92,7 @@ func TestCreateSubnetCommand_Prepare_VPCDeleting(t *testing.T) {
kv.AddInDB(db, "vpc/vpc-dying/state", "deleting")
cmd := CreateSubnetCommand{
Name: "sn-1", VPC: "vpc-dying", VxlanID: 100,
IfaceType: "vms", GatewayIP: "10.0.0.1", CIDR: "10.0.0.0/24",
IfaceType: "vms", InterfaceIP: "10.0.0.1", CIDR: "10.0.0.0/24",
}
if err := cmd.Prepare(db, testCfg()); err == nil {
t.Error("Prepare devrait échouer si le VPC est en cours de suppression")
@ -104,13 +104,109 @@ func TestCreateSubnetCommand_Prepare_VPCDeleted(t *testing.T) {
kv.AddInDB(db, "vpc/vpc-gone/state", "deleted")
cmd := CreateSubnetCommand{
Name: "sn-1", VPC: "vpc-gone", VxlanID: 100,
IfaceType: "vms", GatewayIP: "10.0.0.1", CIDR: "10.0.0.0/24",
IfaceType: "vms", InterfaceIP: "10.0.0.1", CIDR: "10.0.0.0/24",
}
if err := cmd.Prepare(db, testCfg()); err == nil {
t.Error("Prepare devrait échouer si le VPC est supprimé")
}
}
func TestCreateSubnetCommand_Prepare_DefaultsToVxlanMode(t *testing.T) {
_, db := newTestDispatcher(t)
kv.AddInDB(db, "vpc/vpc-1/state", "created")
cmd := CreateSubnetCommand{
Name: "sn-1", VPC: "vpc-1", VxlanID: 100,
IfaceType: "vms", InterfaceIP: "10.0.0.1", CIDR: "10.0.0.0/24",
}
cmd.Prepare(db, testCfg())
mode, _ := kv.GetFromDB(db, "subnet/sn-1/mode")
if mode != "vxlan" {
t.Errorf("mode attendu vxlan, obtenu %q", mode)
}
if _, err := kv.GetFromDB(db, "subnet/sn-1/vxlan_id"); err != nil {
t.Error("vxlan_id devrait être écrit en mode vxlan")
}
}
func TestCreateSubnetCommand_Prepare_BridgeMode_Success(t *testing.T) {
_, db := newTestDispatcher(t)
kv.AddInDB(db, "vpc/vpc-1/state", "created")
cmd := CreateSubnetCommand{
Name: "sn-1", VPC: "vpc-1", Mode: "bridge",
IfaceType: "vms", InterfaceIP: "10.0.0.1", CIDR: "10.0.0.0/24",
}
if err := cmd.Prepare(db, testCfg()); err != nil {
t.Fatalf("Prepare a échoué : %v", err)
}
mode, _ := kv.GetFromDB(db, "subnet/sn-1/mode")
if mode != "bridge" {
t.Errorf("mode attendu bridge, obtenu %q", mode)
}
iface, _ := kv.GetFromDB(db, "subnet/sn-1/local_iface")
if iface != "br-vms" {
t.Errorf("local_iface attendu br-vms, obtenu %q", iface)
}
}
func TestCreateSubnetCommand_Prepare_BridgeMode_NoVxlanID(t *testing.T) {
_, db := newTestDispatcher(t)
kv.AddInDB(db, "vpc/vpc-1/state", "created")
cmd := CreateSubnetCommand{
Name: "sn-1", VPC: "vpc-1", Mode: "bridge",
IfaceType: "vms", InterfaceIP: "10.0.0.1", CIDR: "10.0.0.0/24",
}
cmd.Prepare(db, testCfg())
if _, err := kv.GetFromDB(db, "subnet/sn-1/vxlan_id"); err == nil {
t.Error("vxlan_id ne devrait pas être écrit en mode bridge")
}
}
func TestCreateSubnetCommand_Prepare_UnknownMode(t *testing.T) {
_, db := newTestDispatcher(t)
kv.AddInDB(db, "vpc/vpc-1/state", "created")
cmd := CreateSubnetCommand{
Name: "sn-1", VPC: "vpc-1", Mode: "vlan",
IfaceType: "vms", InterfaceIP: "10.0.0.1", CIDR: "10.0.0.0/24",
}
if err := cmd.Prepare(db, testCfg()); err == nil {
t.Error("Prepare devrait échouer pour un mode inconnu")
}
}
func TestCreateSubnetCommand_Prepare_DefaultRouteStored(t *testing.T) {
_, db := newTestDispatcher(t)
kv.AddInDB(db, "vpc/vpc-1/state", "created")
cmd := CreateSubnetCommand{
Name: "sn-1", VPC: "vpc-1", VxlanID: 100,
IfaceType: "vms", InterfaceIP: "10.0.0.1", CIDR: "10.0.0.0/24",
DefaultRoute: true,
}
if err := cmd.Prepare(db, testCfg()); err != nil {
t.Fatalf("Prepare a échoué : %v", err)
}
val, err := kv.GetFromDB(db, "subnet/sn-1/default_route")
if err != nil {
t.Fatalf("default_route non écrit en DB : %v", err)
}
if val != "true" {
t.Errorf("default_route attendu true, obtenu %q", val)
}
}
func TestCreateSubnetCommand_Prepare_DefaultRouteFalseByDefault(t *testing.T) {
_, db := newTestDispatcher(t)
kv.AddInDB(db, "vpc/vpc-1/state", "created")
cmd := CreateSubnetCommand{
Name: "sn-1", VPC: "vpc-1", VxlanID: 100,
IfaceType: "vms", InterfaceIP: "10.0.0.1", CIDR: "10.0.0.0/24",
}
cmd.Prepare(db, testCfg())
val, _ := kv.GetFromDB(db, "subnet/sn-1/default_route")
if val != "false" {
t.Errorf("default_route attendu false, obtenu %q", val)
}
}
// --- DeleteSubnetCommand.Prepare ---
func TestDeleteSubnetCommand_Prepare_Success(t *testing.T) {

View file

@ -2,7 +2,9 @@ package dispatcher
import (
"fmt"
"math/rand"
"strconv"
"strings"
"time"
configuration "git.g3e.fr/syonad/two/internal/config/agent"
@ -11,16 +13,21 @@ import (
"github.com/dgraph-io/badger/v4"
)
type VMDisk struct {
Path string
Dev string
}
type StartVMCommand struct {
Name string
Subnet string
IP string
MetadataPort string
VolumePath string
Memory int
CPUs int
Password string
SSHKey string
Name string
Subnet string
IP string
Disks []VMDisk
Memory int
CPUs int
UEFI bool
Password string
SSHKey string
}
func (c StartVMCommand) Prepare(db *badger.DB, _ *configuration.Config) error {
@ -34,13 +41,22 @@ func (c StartVMCommand) Prepare(db *badger.DB, _ *configuration.Config) error {
if subnetState == "deleting" || subnetState == "deleted" {
return fmt.Errorf("subnet %q is %s", c.Subnet, subnetState)
}
port, err := allocateMetadataPort(db)
if err != nil {
return fmt.Errorf("allocate metadata port: %w", err)
}
kv.AddInDB(db, "vm/"+c.Name+"/state", "starting")
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", c.MetadataPort)
kv.AddInDB(db, "vm/"+c.Name+"/volume_path", c.VolumePath)
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+"/memory", strconv.Itoa(c.Memory))
kv.AddInDB(db, "vm/"+c.Name+"/cpus", strconv.Itoa(c.CPUs))
if c.UEFI {
kv.AddInDB(db, "vm/"+c.Name+"/uefi", "true")
}
if c.Password != "" {
kv.AddInDB(db, "vm/"+c.Name+"/password", c.Password)
}
@ -50,6 +66,28 @@ func (c StartVMCommand) Prepare(db *badger.DB, _ *configuration.Config) error {
return nil
}
func allocateMetadataPort(db *badger.DB) (int, error) {
entries, err := kv.ListByPrefix(db, "vm/")
if err != nil {
return 0, err
}
used := make(map[int]struct{})
for key, value := range entries {
if strings.HasSuffix(key, "/metadata_port") {
if p, err := strconv.Atoi(value); err == nil {
used[p] = struct{}{}
}
}
}
for range 100 {
p := rand.Intn(9000) + 1000
if _, taken := used[p]; !taken {
return p, nil
}
}
return 0, fmt.Errorf("no free metadata port available in [1000, 9999]")
}
func (c StartVMCommand) Execute(db *badger.DB, cfg *configuration.Config) error {
timeout := time.After(time.Duration(cfg.Dispatcher.TimeoutSeconds) * time.Second)
for {

View file

@ -0,0 +1,130 @@
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à")
}
}

View file

@ -2,6 +2,7 @@ package dispatcher
import (
"fmt"
"net"
"strings"
"time"
@ -13,12 +14,19 @@ import (
type CreateVPCCommand struct {
Name string
CIDR string
}
func (c CreateVPCCommand) Prepare(db *badger.DB, _ *configuration.Config) error {
if _, err := kv.GetFromDB(db, "vpc/"+c.Name+"/state"); err == nil {
return fmt.Errorf("vpc %q already exists", c.Name)
}
if _, _, err := net.ParseCIDR(c.CIDR); err != nil {
return fmt.Errorf("invalid cidr %q: %w", c.CIDR, err)
}
if err := kv.AddInDB(db, "vpc/"+c.Name+"/cidr", c.CIDR); err != nil {
return err
}
return kv.AddInDB(db, "vpc/"+c.Name+"/state", "creating")
}

View file

@ -10,7 +10,7 @@ import (
func TestCreateVPCCommand_Prepare_NewVPC(t *testing.T) {
_, db := newTestDispatcher(t)
cmd := CreateVPCCommand{Name: "vpc-1"}
cmd := CreateVPCCommand{Name: "vpc-1", CIDR: "10.0.0.0/16"}
if err := cmd.Prepare(db, nil); err != nil {
t.Fatalf("Prepare a échoué : %v", err)
}
@ -21,17 +21,32 @@ func TestCreateVPCCommand_Prepare_NewVPC(t *testing.T) {
if state != "creating" {
t.Errorf("state attendu creating, obtenu %q", state)
}
cidr, err := kv.GetFromDB(db, "vpc/vpc-1/cidr")
if err != nil {
t.Fatalf("cidr non écrit en DB : %v", err)
}
if cidr != "10.0.0.0/16" {
t.Errorf("cidr attendu 10.0.0.0/16, obtenu %q", cidr)
}
}
func TestCreateVPCCommand_Prepare_Duplicate(t *testing.T) {
_, db := newTestDispatcher(t)
kv.AddInDB(db, "vpc/vpc-exist/state", "created")
cmd := CreateVPCCommand{Name: "vpc-exist"}
cmd := CreateVPCCommand{Name: "vpc-exist", CIDR: "10.0.0.0/16"}
if err := cmd.Prepare(db, nil); err == nil {
t.Error("Prepare devrait échouer sur un VPC déjà existant")
}
}
func TestCreateVPCCommand_Prepare_InvalidCIDR(t *testing.T) {
_, db := newTestDispatcher(t)
cmd := CreateVPCCommand{Name: "vpc-bad", CIDR: "not-a-cidr"}
if err := cmd.Prepare(db, nil); err == nil {
t.Error("Prepare devrait échouer avec un CIDR invalide")
}
}
// --- DeleteVPCCommand.Prepare ---
func TestDeleteVPCCommand_Prepare_Success(t *testing.T) {

View file

@ -13,46 +13,48 @@ func deleteRule(args ...string) error {
return exec.Command("ebtables", append([]string{"-D"}, args...)...).Run()
}
func DropARPToGateway(bridge, gatewayIP string) error {
func DropARPToGateway(iface, ip string) error {
if err := addRule("FORWARD",
"--out-interface", bridge,
"--out-interface", iface,
"-p", "arp",
"--arp-op", "Request",
"--arp-ip-dst", gatewayIP,
"--arp-ip-dst", ip,
"-j", "DROP"); err != nil {
return fmt.Errorf("ebtables arp rule: %w", err)
}
return nil
}
func DropDHCP(bridge string) error {
func DropDHCP(iface, ip string) error {
if err := addRule("FORWARD",
"--out-interface", bridge,
"--out-interface", iface,
"-p", "IPv4",
"--ip-protocol", "udp",
"--ip-source-port", "67:68",
"--ip-destination-port", "67:68",
"--ip-source", ip,
"-j", "DROP"); err != nil {
return fmt.Errorf("ebtables dhcp rule: %w", err)
}
return nil
}
func DeleteARPToGateway(bridge, gatewayIP string) error {
func DeleteARPToGateway(iface, ip string) error {
return deleteRule("FORWARD",
"--out-interface", bridge,
"--out-interface", iface,
"-p", "arp",
"--arp-op", "Request",
"--arp-ip-dst", gatewayIP,
"--arp-ip-dst", ip,
"-j", "DROP")
}
func DeleteDHCP(bridge string) error {
func DeleteDHCP(iface, ip string) error {
return deleteRule("FORWARD",
"--out-interface", bridge,
"--out-interface", iface,
"-p", "IPv4",
"--ip-protocol", "udp",
"--ip-source-port", "67:68",
"--ip-destination-port", "67:68",
"--ip-source", ip,
"-j", "DROP")
}

View file

@ -2,12 +2,8 @@
users:
- name: syonad
lock_passwd: false
gecos: alpine Cloud User
groups: [adm, wheel]
doas:
- permit nopass syonad
sudo: ["ALL=(ALL) NOPASSWD:ALL"]
shell: /bin/ash
shell: /bin/bash
passwd: "{{ .Password }}"
ssh_authorized_keys:
- "{{ .SSHKEY }}"
- "{{ .SSHKEY }}"

View file

@ -0,0 +1,29 @@
//go:build linux
package netif
import (
"fmt"
"net"
"github.com/vishvananda/netlink"
)
func GetDefaultGateway() (net.IP, error) {
routes, err := netlink.RouteList(nil, netlink.FAMILY_V4)
if err != nil {
return nil, fmt.Errorf("list routes: %w", err)
}
for _, r := range routes {
if r.Gw == nil {
continue
}
if r.Dst == nil {
return r.Gw, nil
}
if ones, _ := r.Dst.Mask.Size(); ones == 0 {
return r.Gw, nil
}
}
return nil, fmt.Errorf("no default gateway found")
}

View file

@ -0,0 +1,12 @@
//go:build !linux
package netif
import (
"fmt"
"net"
)
func GetDefaultGateway() (net.IP, error) {
return nil, fmt.Errorf("not supported on this platform")
}

20
internal/qemu/config.go Normal file
View file

@ -0,0 +1,20 @@
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
}

View file

@ -4,18 +4,13 @@ package qemu
import (
"fmt"
"os"
"os/exec"
"path/filepath"
"sort"
"strings"
)
type Config struct {
Name string
TapID int
Mac string
VolumePath string
Memory int
CPUs int
}
func Start(cfg Config) error {
memory := cfg.Memory
if memory == 0 {
@ -27,20 +22,83 @@ func Start(cfg Config) error {
cpus = 1
}
cmd := exec.Command("qemu-system-x86_64",
for _, dir := range []string{cfg.SerialDir, cfg.MonitorDir, cfg.QMPDir} {
if dir != "" {
if err := os.MkdirAll(dir, 0755); err != nil {
return fmt.Errorf("mkdir %s: %w", dir, err)
}
}
}
serialSock := filepath.Join(cfg.SerialDir, cfg.Name+".sock")
monitorSock := filepath.Join(cfg.MonitorDir, cfg.Name+".sock")
qmpSock := filepath.Join(cfg.QMPDir, cfg.Name+".sock")
args := []string{
"-enable-kvm",
"-cpu", "host",
"-m", fmt.Sprintf("%d", memory),
"-smp", fmt.Sprintf("%d", cpus),
"-serial", fmt.Sprintf("unix:/tmp/%s.sock,server,nowait", cfg.Name),
"-monitor", fmt.Sprintf("unix:/tmp/%s.mon-sock,server,nowait", cfg.Name),
"-qmp", fmt.Sprintf("unix:/tmp/%s.qmp-sock,server,nowait", cfg.Name),
"-serial", fmt.Sprintf("unix:%s,server,nowait", serialSock),
"-monitor", fmt.Sprintf("unix:%s,server,nowait", monitorSock),
"-qmp", fmt.Sprintf("unix:%s,server,nowait", qmpSock),
"-display", "none",
"-drive", fmt.Sprintf("file=%s,if=virtio", cfg.VolumePath),
}
if cfg.UEFICodePath != "" && cfg.UEFIVarsPath != "" {
args = append(args,
"-drive", fmt.Sprintf("if=pflash,format=raw,readonly=on,file=%s", cfg.UEFICodePath),
"-drive", fmt.Sprintf("if=pflash,format=raw,file=%s", cfg.UEFIVarsPath),
)
}
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,
"-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",
)
cmd := exec.Command("qemu-system-x86_64", args...)
if err := cmd.Run(); err != nil {
return fmt.Errorf("qemu-system-x86_64: %w", err)
}

View file

@ -2,12 +2,9 @@
package qemu
import "errors"
type Config struct {
Name, Mac, VolumePath string
TapID, Memory, CPUs int
}
import (
"errors"
)
func Start(_ Config) error {
return errors.New("vm: not supported on this platform")

View file

@ -27,46 +27,48 @@ func CreateSubnet(db *badger.DB, subnetName string) error {
return err
}
vxlanIface := fmt.Sprintf("vxlan-%d", d.vxlanID)
if err := createSubnet(db, subnetName, d); err != nil {
return err
}
if err := netif.CreateVethToNetns("v-"+d.subnetID+"-e", "v-"+d.subnetID+"-i", "/var/run/netns/"+d.vpc, 1500); err != nil {
return kv.AddInDB(db, "subnet/"+subnetName+"/state", "created")
}
func createSubnet(db *badger.DB, subnetName string, d subnetData) error {
vethE := "v-" + d.subnetID + "-e"
vethI := "v-" + d.subnetID + "-i"
if err := netif.CreateVethToNetns(vethE, vethI, "/var/run/netns/"+d.vpc, 1500); err != nil {
return fmt.Errorf("create veth: %w", err)
}
if err := netif.CreateBridge(d.bridge, 1500); err != nil {
return fmt.Errorf("create bridge: %w", err)
}
if err := netns.Call(d.vpc, func() error {
return netif.CreateBridge(d.bridge, 1500)
}); err != nil {
return fmt.Errorf("create bridge in netns: %w", err)
}
if err := netif.CreateVxlan(vxlanIface, d.vxlanID, d.localIface, 1500); err != nil {
return fmt.Errorf("create vxlan: %w", err)
}
if err := netif.BridgeSetMaster("v-"+d.subnetID+"-e", d.bridge); err != nil {
return fmt.Errorf("add veth-e to bridge: %w", err)
}
if err := netns.Call(d.vpc, func() error {
return netif.BridgeSetMaster("v-"+d.subnetID+"-i", d.bridge)
}); err != nil {
return fmt.Errorf("add veth-i to bridge in netns: %w", err)
}
if err := netif.BridgeSetMaster(vxlanIface, d.bridge); err != nil {
return fmt.Errorf("add vxlan to bridge: %w", err)
}
for _, iface := range []string{"v-" + d.subnetID + "-e", vxlanIface, d.bridge} {
if err := netif.LinkSetUp(iface); err != nil {
return fmt.Errorf("set up %s: %w", iface, err)
switch d.mode {
case "vxlan":
if err := setupVxlanHost(d, vethE); err != nil {
return err
}
case "bridge":
if err := netif.BridgeSetMaster(vethE, d.localIface); err != nil {
return fmt.Errorf("add veth-e to bridge: %w", err)
}
if err := netif.LinkSetUp(vethE); err != nil {
return fmt.Errorf("set up %s: %w", vethE, err)
}
default:
return fmt.Errorf("unknown subnet mode %q", d.mode)
}
if err := netns.Call(d.vpc, func() error {
for _, iface := range []string{"v-" + d.subnetID + "-i", d.bridge} {
if err := netif.CreateBridge(d.bridge, 1500); err != nil {
return fmt.Errorf("create bridge: %w", err)
}
return netif.BridgeSetMaster(vethI, d.bridge)
}); err != nil {
return fmt.Errorf("setup bridge in netns: %w", err)
}
if err := netns.Call(d.vpc, func() error {
for _, iface := range []string{vethI, d.bridge} {
if err := netif.LinkSetUp(iface); err != nil {
return fmt.Errorf("set up %s: %w", iface, err)
}
@ -76,31 +78,81 @@ func CreateSubnet(db *badger.DB, subnetName string) error {
return fmt.Errorf("set up interfaces in netns: %w", err)
}
if err := netns.Call(d.vpc, func() error {
return netif.AddrAdd(d.bridge, d.gatewayIP)
}); err != nil {
return fmt.Errorf("add addr to bridge in netns: %w", err)
switch d.mode {
case "vxlan":
if err := netns.Call(d.vpc, func() error {
if err := netif.AddrAdd(d.bridge, d.interfaceIP); err != nil {
return fmt.Errorf("add addr: %w", err)
}
if err := netif.RouteAdd(d.bridge, d.cidr); err != nil {
return fmt.Errorf("add route: %w", err)
}
if err := ebtables.DropARPToGateway(vethI, d.interfaceIP.String()); err != nil {
return err
}
return ebtables.DropDHCP(vethI, d.interfaceIP.String())
}); err != nil {
return fmt.Errorf("configure netns: %w", err)
}
case "bridge":
if err := netns.Call(d.vpc, func() error {
if err := netif.AddrAdd(d.bridge, d.interfaceIP); err != nil {
return fmt.Errorf("add addr: %w", err)
}
if err := netif.RouteAdd(d.bridge, d.cidr); err != nil {
return fmt.Errorf("add route: %w", err)
}
return ebtables.DropDHCP(vethI, d.interfaceIP.String())
}); err != nil {
return fmt.Errorf("configure netns: %w", err)
}
}
if err := netns.Call(d.vpc, func() error {
return netif.RouteAdd(d.bridge, d.cidr)
}); err != nil {
return fmt.Errorf("add route in netns: %w", err)
}
return startDHCP(db, subnetName, d)
}
if err := ebtables.DropARPToGateway(d.bridge, d.gatewayIP.String()); err != nil {
return err
}
if err := ebtables.DropDHCP(d.bridge); err != nil {
return err
}
func setupVxlanHost(d subnetData, vethE string) error {
vxlanIface := fmt.Sprintf("vxlan-%d", d.vxlanID)
if err := netif.CreateBridge(d.bridge, 1500); err != nil {
return fmt.Errorf("create bridge: %w", err)
}
if err := netif.CreateVxlan(vxlanIface, d.vxlanID, d.localIface, 1500); err != nil {
return fmt.Errorf("create vxlan: %w", err)
}
if err := netif.BridgeSetMaster(vethE, d.bridge); err != nil {
return fmt.Errorf("add veth-e to bridge: %w", err)
}
if err := netif.BridgeSetMaster(vxlanIface, d.bridge); err != nil {
return fmt.Errorf("add vxlan to bridge: %w", err)
}
for _, iface := range []string{vethE, vxlanIface, d.bridge} {
if err := netif.LinkSetUp(iface); err != nil {
return fmt.Errorf("set up %s: %w", iface, err)
}
}
return nil
}
func startDHCP(db *badger.DB, subnetName string, d subnetData) error {
conf := dhcp.Config{
Network: d.cidr,
Gateway: d.gatewayIP,
Name: d.vpc + "_" + d.bridge,
ConfDir: "/etc/dnsmasq.d",
}
switch d.mode {
case "vxlan":
conf.VPCGateway = d.interfaceIP
conf.VPCRoute = d.vpcCIDR
case "bridge":
if d.defaultRoute {
gw, err := netif.GetDefaultGateway()
if err != nil {
return fmt.Errorf("get default gateway: %w", err)
}
conf.DefaultGateway = gw
}
}
_, entries, err := dhcp.GenerateConfig(conf)
if err != nil {
return fmt.Errorf("generate dhcp config: %w", err)
@ -118,6 +170,5 @@ func CreateSubnet(db *badger.DB, subnetName string) error {
if err := svc.Start("dnsmasq@" + conf.Name + ".service"); err != nil {
return fmt.Errorf("start dnsmasq: %w", err)
}
return kv.AddInDB(db, "subnet/"+subnetName+"/state", "created")
return nil
}

View file

@ -11,13 +11,16 @@ import (
)
type subnetData struct {
vpc string
subnetID string
bridge string
vxlanID int
localIface string
gatewayIP net.IP
cidr *net.IPNet
vpc string
subnetID string
bridge string
mode string
vxlanID int
localIface string
interfaceIP net.IP
cidr *net.IPNet
vpcCIDR *net.IPNet
defaultRoute bool
}
func loadSubnet(db *badger.DB, name string) (subnetData, error) {
@ -32,15 +35,23 @@ func loadSubnet(db *badger.DB, name string) (subnetData, error) {
}
d.vpc = vpc
vxlanIDStr, err := kv.GetFromDB(db, "subnet/"+name+"/vxlan_id")
mode, err := kv.GetFromDB(db, "subnet/"+name+"/mode")
if err != nil {
return d, fmt.Errorf("get vxlan_id: %w", err)
return d, fmt.Errorf("get mode: %w", err)
}
vxlanID, err := strconv.Atoi(vxlanIDStr)
if err != nil {
return d, fmt.Errorf("parse vxlan_id: %w", err)
d.mode = mode
if d.mode == "vxlan" {
vxlanIDStr, err := kv.GetFromDB(db, "subnet/"+name+"/vxlan_id")
if err != nil {
return d, fmt.Errorf("get vxlan_id: %w", err)
}
vxlanID, err := strconv.Atoi(vxlanIDStr)
if err != nil {
return d, fmt.Errorf("parse vxlan_id: %w", err)
}
d.vxlanID = vxlanID
}
d.vxlanID = vxlanID
localIface, err := kv.GetFromDB(db, "subnet/"+name+"/local_iface")
if err != nil {
@ -48,15 +59,15 @@ func loadSubnet(db *badger.DB, name string) (subnetData, error) {
}
d.localIface = localIface
gatewayIPStr, err := kv.GetFromDB(db, "subnet/"+name+"/gateway_ip")
interfaceIPStr, err := kv.GetFromDB(db, "subnet/"+name+"/interface_ip")
if err != nil {
return d, fmt.Errorf("get gateway_ip: %w", err)
return d, fmt.Errorf("get interface_ip: %w", err)
}
gatewayIP := net.ParseIP(gatewayIPStr)
if gatewayIP == nil {
return d, fmt.Errorf("invalid gateway_ip: %s", gatewayIPStr)
interfaceIP := net.ParseIP(interfaceIPStr)
if interfaceIP == nil {
return d, fmt.Errorf("invalid interface_ip: %s", interfaceIPStr)
}
d.gatewayIP = gatewayIP
d.interfaceIP = interfaceIP
cidrStr, err := kv.GetFromDB(db, "subnet/"+name+"/cidr")
if err != nil {
@ -68,5 +79,21 @@ func loadSubnet(db *badger.DB, name string) (subnetData, error) {
}
d.cidr = ipNet
defaultRouteStr, err := kv.GetFromDB(db, "subnet/"+name+"/default_route")
if err != nil {
return d, fmt.Errorf("get default_route: %w", err)
}
d.defaultRoute = defaultRouteStr == "true"
vpcCIDRStr, err := kv.GetFromDB(db, "vpc/"+d.vpc+"/cidr")
if err != nil {
return d, fmt.Errorf("get vpc cidr: %w", err)
}
_, vpcIPNet, err := net.ParseCIDR(vpcCIDRStr)
if err != nil {
return d, fmt.Errorf("parse vpc cidr: %w", err)
}
d.vpcCIDR = vpcIPNet
return d, nil
}

View file

@ -27,36 +27,64 @@ func DeleteSubnet(db *badger.DB, subnetName string) error {
return err
}
vxlanIface := fmt.Sprintf("vxlan-%d", d.vxlanID)
if err := stopDHCP(db, subnetName, d); err != nil {
return err
}
switch d.mode {
case "vxlan":
if err := deleteSubnetVxlan(d); err != nil {
return err
}
case "bridge":
if err := deleteSubnetBridge(d); err != nil {
return err
}
default:
return fmt.Errorf("unknown subnet mode %q", d.mode)
}
return kv.AddInDB(db, "subnet/"+subnetName+"/state", "deleted")
}
func stopDHCP(db *badger.DB, subnetName string, d subnetData) error {
svc, err := systemd.New()
if err != nil {
return fmt.Errorf("connect to systemd: %w", err)
}
defer svc.Close()
if err := svc.Stop("dnsmasq@" + d.vpc + "_" + d.bridge + ".service"); err != nil {
return fmt.Errorf("stop dnsmasq: %w", err)
svcName := "dnsmasq@" + d.vpc + "_" + d.bridge + ".service"
if status, err := svc.Status(svcName); err == nil && status.ActiveState == "active" {
if err := svc.Stop(svcName); err != nil {
return fmt.Errorf("stop dnsmasq: %w", err)
}
}
if err := os.Remove("/etc/dnsmasq.d/" + d.vpc + "_" + d.bridge + ".conf"); err != nil && !os.IsNotExist(err) {
return fmt.Errorf("remove dnsmasq config: %w", err)
}
if err := kv.DeleteInDB(db, "subnet/"+subnetName+"/dhcp"); err != nil {
return fmt.Errorf("delete dhcp entries: %w", err)
}
return nil
}
if err := ebtables.DeleteARPToGateway(d.bridge, d.gatewayIP.String()); err != nil {
return fmt.Errorf("delete ebtables arp rule: %w", err)
}
if err := ebtables.DeleteDHCP(d.bridge); err != nil {
return fmt.Errorf("delete ebtables dhcp rule: %w", err)
}
func deleteSubnetVxlan(d subnetData) error {
vxlanIface := fmt.Sprintf("vxlan-%d", d.vxlanID)
vethI := "v-" + d.subnetID + "-i"
if err := netns.Call(d.vpc, func() error {
if err := ebtables.DeleteARPToGateway(vethI, d.interfaceIP.String()); err != nil {
return fmt.Errorf("delete ebtables arp rule: %w", err)
}
if err := ebtables.DeleteDHCP(vethI, d.interfaceIP.String()); err != nil {
return fmt.Errorf("delete ebtables dhcp rule: %w", err)
}
return netif.DeleteLink(d.bridge)
}); err != nil {
return fmt.Errorf("delete bridge in netns: %w", err)
return fmt.Errorf("delete netns resources: %w", err)
}
if err := netif.DeleteLink(vxlanIface); err != nil {
@ -70,6 +98,23 @@ func DeleteSubnet(db *badger.DB, subnetName string) error {
if err := netif.DeleteLink(d.bridge); err != nil {
return fmt.Errorf("delete bridge: %w", err)
}
return kv.AddInDB(db, "subnet/"+subnetName+"/state", "deleted")
return nil
}
func deleteSubnetBridge(d subnetData) error {
vethI := "v-" + d.subnetID + "-i"
if err := netns.Call(d.vpc, func() error {
if err := ebtables.DeleteDHCP(vethI, d.interfaceIP.String()); err != nil {
return fmt.Errorf("delete ebtables dhcp rule: %w", err)
}
return netif.DeleteLink(d.bridge)
}); err != nil {
return fmt.Errorf("delete netns resources: %w", err)
}
if err := netif.DeleteLink("v-" + d.subnetID + "-e"); err != nil {
return fmt.Errorf("delete veth: %w", err)
}
return nil
}

View file

@ -2,6 +2,9 @@ package vm
import (
"fmt"
"io"
"os"
"path/filepath"
configuration "git.g3e.fr/syonad/two/internal/config/agent"
"git.g3e.fr/syonad/two/internal/iptables"
@ -33,7 +36,7 @@ func StartVM(db *badger.DB, name string, cfg *configuration.Config) error {
}
if err := netns.Call(d.vpcName, func() error {
return iptables.AddMetadataRedirect(d.ip, d.gatewayIP, d.metadataPort)
return iptables.AddMetadataRedirect(d.ip, d.interfaceIP, d.metadataPort)
}); err != nil {
return fmt.Errorf("add metadata redirect: %w", err)
}
@ -41,7 +44,7 @@ func StartVM(db *badger.DB, name string, cfg *configuration.Config) error {
if err := metadata.StartMetadata(metadata.NoCloudConfig{
Name: name,
VpcName: d.vpcName,
BindIP: d.gatewayIP,
BindIP: d.interfaceIP,
BindPort: d.metadataPort,
Password: d.password,
SSHKEY: d.sshkey,
@ -49,18 +52,59 @@ 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,
Memory: d.memory,
CPUs: d.cpus,
SerialDir: cfg.QEMU.SerialDir,
MonitorDir: cfg.QEMU.MonitorDir,
QMPDir: cfg.QEMU.QMPDir,
}
if d.uefi {
varsPath := filepath.Join(cfg.QEMU.UEFIVarsDir, name+"-uefi-vars.fd")
if err := copyFile(cfg.QEMU.OVMFVarsTemplate, varsPath); err != nil {
return fmt.Errorf("copy uefi vars: %w", err)
}
qcfg.UEFICodePath = cfg.QEMU.OVMFCodePath
qcfg.UEFIVarsPath = varsPath
}
if err := netns.Call(d.vpcName, func() error {
return qemu.Start(qemu.Config{
Name: name,
TapID: d.tapID,
Mac: d.mac,
VolumePath: d.volumePath,
Memory: d.memory,
CPUs: d.cpus,
})
return qemu.Start(qcfg)
}); err != nil {
return fmt.Errorf("start qemu: %w", err)
}
return kv.AddInDB(db, "vm/"+name+"/state", "started")
}
func copyFile(src, dst string) error {
if err := os.MkdirAll(filepath.Dir(dst), 0755); err != nil {
return err
}
in, err := os.Open(src)
if err != nil {
return err
}
defer in.Close()
out, err := os.Create(dst)
if err != nil {
return err
}
defer out.Close()
if _, err := io.Copy(out, in); err != nil {
return err
}
return out.Sync()
}

View file

@ -11,18 +11,24 @@ import (
"github.com/dgraph-io/badger/v4"
)
type diskEntry struct {
path string
dev string
}
type vmData struct {
subnetName string
vpcName string
gatewayIP string
interfaceIP string
bridge string
tapID int
ip string
metadataPort string
mac string
volumePath string
disks []diskEntry
memory int
cpus int
uefi bool
password string
sshkey string
}
@ -43,11 +49,11 @@ func loadVM(db *badger.DB, name string) (vmData, error) {
}
d.vpcName = vpcName
gatewayIP, err := kv.GetFromDB(db, "subnet/"+subnetName+"/gateway_ip")
interfaceIP, err := kv.GetFromDB(db, "subnet/"+subnetName+"/interface_ip")
if err != nil {
return d, fmt.Errorf("get gateway_ip: %w", err)
return d, fmt.Errorf("get interface_ip: %w", err)
}
d.gatewayIP = gatewayIP
d.interfaceIP = interfaceIP
tapIDStr, err := kv.GetFromDB(db, "vm/"+name+"/tap_id")
if err != nil {
@ -81,11 +87,18 @@ func loadVM(db *badger.DB, name string) (vmData, error) {
}
d.mac = mac
volumePath, err := kv.GetFromDB(db, "vm/"+name+"/volume_path")
diskEntries, err := kv.ListByPrefix(db, "vm/"+name+"/disk/")
if err != nil {
return d, fmt.Errorf("get volume_path: %w", err)
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})
}
d.volumePath = volumePath
memoryStr, err := kv.GetFromDB(db, "vm/"+name+"/memory")
if err != nil {
@ -105,6 +118,9 @@ func loadVM(db *badger.DB, name string) (vmData, error) {
return d, fmt.Errorf("parse cpus: %w", err)
}
if v, _ := kv.GetFromDB(db, "vm/"+name+"/uefi"); v == "true" {
d.uefi = true
}
d.password, _ = kv.GetFromDB(db, "vm/"+name+"/password")
d.sshkey, _ = kv.GetFromDB(db, "vm/"+name+"/sshkey")

View file

@ -2,6 +2,8 @@ package vm
import (
"fmt"
"os"
"path/filepath"
"time"
configuration "git.g3e.fr/syonad/two/internal/config/agent"
@ -29,30 +31,23 @@ func StopVM(db *badger.DB, name string, cfg *configuration.Config) error {
return err
}
socketPath := fmt.Sprintf("/tmp/%s.qmp-sock", name)
socketPath := filepath.Join(cfg.QEMU.QMPDir, name+".sock")
if _, err := qmp.Send(socketPath, []string{`{"execute":"system_powerdown"}`}); err != nil {
return fmt.Errorf("qmp system_powerdown: %w", err)
}
// 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 := 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
}
// socket absent ou QEMU déjà arrêté : cleanup direct
if err := netns.Call(d.vpcName, func() error {
return iptables.DeleteMetadataRedirect(d.ip, d.gatewayIP, d.metadataPort)
return iptables.DeleteMetadataRedirect(d.ip, d.interfaceIP, d.metadataPort)
}); err != nil {
return fmt.Errorf("delete metadata redirect: %w", err)
}
@ -65,5 +60,25 @@ func StopVM(db *badger.DB, name string, cfg *configuration.Config) error {
return fmt.Errorf("delete tap: %w", err)
}
if d.uefi {
varsPath := filepath.Join(cfg.QEMU.UEFIVarsDir, name+"-uefi-vars.fd")
os.Remove(varsPath)
}
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
}
}
}
}

52
pkg/db/kv/admin_server.go Normal file
View file

@ -0,0 +1,52 @@
package kv
import (
"fmt"
"log/slog"
"net/http"
"sort"
"github.com/dgraph-io/badger/v4"
)
type AdminServer struct {
db *badger.DB
logger *slog.Logger
}
func NewAdminServer(db *badger.DB, logger *slog.Logger) *AdminServer {
return &AdminServer{db: db, logger: logger}
}
func (s *AdminServer) Start(address string) {
mux := http.NewServeMux()
mux.HandleFunc("/db", s.dbHandler)
s.logger.Info("admin server listening", "address", address)
if err := http.ListenAndServe(address, mux); err != nil {
s.logger.Error("admin server stopped", "error", err)
}
}
func (s *AdminServer) dbHandler(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodGet {
http.Error(w, "method not allowed", http.StatusMethodNotAllowed)
return
}
entries, err := ListByPrefix(s.db, r.URL.Query().Get("prefix"))
if err != nil {
http.Error(w, "db error: "+err.Error(), http.StatusInternalServerError)
return
}
keys := make([]string, 0, len(entries))
for k := range entries {
keys = append(keys, k)
}
sort.Strings(keys)
w.Header().Set("Content-Type", "text/plain; charset=utf-8")
for _, k := range keys {
fmt.Fprintf(w, "%s=%s\n", k, entries[k])
}
}

View file

@ -17,7 +17,8 @@ func InitDB(conf Config, readonly bool) *badger.DB {
opts.NumLevelZeroTablesStall = 2
db, err := badger.Open(opts)
if err != nil {
log.Fatalf("kv.InitDB (readonly=%v, path=%s): %v", readonly, conf.Path, err)
log.Printf("kv.InitDB (readonly=%v, path=%s): %v", readonly, conf.Path, err)
panic(err)
}
return db
}