Compare commits

..

34 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
87808312f4
f-25: bin: remove old binaries
All checks were successful
Pre Release Workflow / set-release-target (push) Successful in 1s
Pre Release Workflow / build (metadata, amd64, linux) (push) Successful in 1m32s
Pre Release Workflow / prerelease (push) Successful in 10s
Pre Release Workflow / build (agent, amd64, linux) (push) Successful in 1m35s
Pre Release Workflow / upload-scripts (run-dnsmasq-in-netns.sh) (push) Successful in 6s
Signed-off-by: GnomeZworc <nicolas.boufidjeline@g3e.fr>
2026-04-29 00:50:04 +02:00
862406f041
f-25: use file for metadata
Signed-off-by: GnomeZworc <nicolas.boufidjeline@g3e.fr>
2026-04-29 00:48:55 +02:00
f5707c343c
f-25: add load files
Signed-off-by: GnomeZworc <nicolas.boufidjeline@g3e.fr>
2026-04-29 00:41:43 +02:00
915581904c
f-25: code: add cfg to start vm
Signed-off-by: GnomeZworc <nicolas.boufidjeline@g3e.fr>
2026-04-29 00:36:01 +02:00
4bef0f9d5f
f-25: config: add metadata dir
Signed-off-by: GnomeZworc <nicolas.boufidjeline@g3e.fr>
2026-04-29 00:33:22 +02:00
e545b70d53
f-25: add propre log for init db
Signed-off-by: GnomeZworc <nicolas.boufidjeline@g3e.fr>
2026-04-29 00:21:55 +02:00
46 changed files with 1438 additions and 482 deletions

View file

@ -33,9 +33,7 @@ jobs:
goos: [linux] goos: [linux]
goarch: [amd64] goarch: [amd64]
binaries: binaries:
- db
- metadata - metadata
- metacli
- agent - agent
uses: ./.forgejo/workflows/build.yml uses: ./.forgejo/workflows/build.yml
with: with:

View file

@ -294,13 +294,17 @@ components:
VPCCreateRequest: VPCCreateRequest:
type: object type: object
required: [name] required: [name, cidr]
properties: properties:
name: name:
type: string type: string
description: Unique name for the VPC, must follow the format vp-[id] description: Unique name for the VPC, must follow the format vp-[id]
pattern: '^vp-.+' pattern: '^vp-.+'
example: vp-00001 example: vp-00001
cidr:
type: string
description: CIDR block for the entire VPC address space
example: "10.0.0.0/16"
VPC: VPC:
type: object type: object
@ -312,10 +316,13 @@ components:
type: string type: string
enum: [creating, created, deleting, deleted] enum: [creating, created, deleting, deleted]
example: created example: created
cidr:
type: string
example: "10.0.0.0/16"
SubnetCreateRequest: SubnetCreateRequest:
type: object type: object
required: [name, vpc, vxlan_id, gateway_ip, cidr] required: [name, vpc, interface_ip, cidr]
properties: properties:
name: name:
type: string type: string
@ -325,15 +332,24 @@ components:
type: string type: string
description: Parent VPC name description: Parent VPC name
example: vpc1 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: vxlan_id:
type: integer type: integer
description: VXLAN VNI identifier description: VXLAN VNI identifier. Required when mode is "vxlan", ignored otherwise.
example: 100 example: 100
iface_type: iface_type:
type: string 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. 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 example: vms
gateway_ip: interface_ip:
type: string type: string
format: ipv4 format: ipv4
description: Gateway IP for the subnet description: Gateway IP for the subnet
@ -342,6 +358,12 @@ components:
type: string type: string
description: Subnet CIDR block description: Subnet CIDR block
example: "10.10.10.0/24" 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: Subnet:
type: object type: object
@ -356,30 +378,35 @@ components:
vpc: vpc:
type: string type: string
example: vpc1 example: vpc1
mode:
type: string
enum: [vxlan, bridge]
example: vxlan
vxlan_id: vxlan_id:
type: integer type: integer
description: VXLAN VNI. Present only when mode is "vxlan".
example: 100 example: 100
local_iface: local_iface:
type: string type: string
description: Resolved interface name description: Resolved interface name from agent config
example: br-000000 example: br-000000
gateway_ip: interface_ip:
type: string type: string
example: "10.10.10.1" example: "10.10.10.1"
cidr: cidr:
type: string type: string
example: "10.10.10.0/24" example: "10.10.10.0/24"
default_route:
type: boolean
example: false
VMCreateRequest: VMCreateRequest:
type: object type: object
required: [name, metadata_port, interfaces, storage] required: [name, interfaces, storage]
properties: properties:
name: name:
type: string type: string
example: vm-00001 example: vm-00001
metadata_port:
type: string
example: "80"
memory: memory:
type: integer type: integer
description: Memory in MB (default 512) description: Memory in MB (default 512)
@ -403,6 +430,10 @@ components:
minItems: 1 minItems: 1
items: items:
$ref: "#/components/schemas/VMStorage" $ref: "#/components/schemas/VMStorage"
uefi:
type: boolean
description: Boot with UEFI firmware (OVMF). Defaults to false (SeaBIOS).
example: false
VMInterface: VMInterface:
type: object type: object
@ -460,6 +491,9 @@ components:
type: array type: array
items: items:
$ref: "#/components/schemas/VMStorage" $ref: "#/components/schemas/VMStorage"
uefi:
type: boolean
example: false
Error: Error:
type: object type: object

View file

@ -51,6 +51,10 @@ func main() {
d := dispatcher.New(q, db, cfg, log.With(slog.String("component", "dispatcher"))) 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 agentapi.New(d, db, log.With(slog.String("component", "api"))).Start(apiAddr)
go promserver.Start(promAddr, registry) 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 {} select {}
} }

View file

@ -1,53 +0,0 @@
package main
import (
"flag"
"fmt"
configuration "git.g3e.fr/syonad/two/internal/config/agent"
"git.g3e.fr/syonad/two/internal/metadata"
"git.g3e.fr/syonad/two/pkg/db/kv"
)
func main() {
conf_file := flag.String("conf", "/etc/two/agent.yml", "configuration file")
vm_name := flag.String("vm_name", "", "Nom de la vm")
vpc := flag.String("vpc_name", "", "vpc name")
bind_ip := flag.String("ip", "", "bind ip")
bind_port := flag.String("port", "", "bind port")
ssh_key := flag.String("key", "", "Clef ssh")
password := flag.String("pass", "", "password user")
start := flag.Bool("start", false, "start metadata server")
stop := flag.Bool("stop", false, "stop metadata server")
dryrun := flag.Bool("dryrun", false, "launch in dry node")
flag.Parse()
conf, err := configuration.LoadConfig(*conf_file)
if err != nil {
fmt.Println(err)
return
}
db := kv.InitDB(kv.Config{
Path: conf.Database.Path,
}, false)
defer db.Close()
if *start {
if err := metadata.StartMetadata(metadata.NoCloudConfig{
VpcName: *vpc,
Name: *vm_name,
BindIP: *bind_ip,
BindPort: *bind_port,
Password: *password,
SSHKEY: *ssh_key,
}, db, *dryrun); err != nil {
fmt.Println(err)
}
} else if *stop {
if err := metadata.StopMetadata(*vm_name, db, *dryrun); err != nil {
fmt.Println(err)
}
}
}

View file

@ -2,26 +2,29 @@ package main
import ( import (
"flag" "flag"
"fmt"
"os"
configuration "git.g3e.fr/syonad/two/internal/config/agent"
"git.g3e.fr/syonad/two/internal/metadata" "git.g3e.fr/syonad/two/internal/metadata"
) )
var ( var (
iface = flag.String("interface", "0.0.0.0", "Interface IP à écouter") confFile = flag.String("conf", "/etc/two/agent.yml", "configuration file")
port = flag.Int("port", 0, "Port à utiliser") vm_name = flag.String("vm", "", "Name of the vm")
netns_name = flag.String("netns", "", "Network namespace à utiliser")
conf_file = flag.String("conf", "/etc/two/agent.yml", "configuration file")
vm_name = flag.String("vm", "", "Name of the vm")
) )
func main() { func main() {
flag.Parse() flag.Parse()
cfg, err := configuration.LoadConfig(*confFile)
if err != nil {
fmt.Fprintf(os.Stderr, "failed to load config: %v\n", err)
os.Exit(1)
}
metadata.StartServer(metadata.ServerConfig{ metadata.StartServer(metadata.ServerConfig{
Netns: *netns_name, VmName: *vm_name,
Iface: *iface, RunDir: cfg.Metadata.RunDir,
Port: *port,
ConfFile: *conf_file,
VmName: *vm_name,
}) })
} }

View file

@ -35,6 +35,28 @@ interfaces:
internet: br-000000 internet: br-000000
admin: br-000000 admin: br-000000
# Metadata server runtime directory (cloud-init files per VM)
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 # Logging configuration
logger: logger:
# Log level: debug, info, warn, error (default: info) # Log level: debug, info, warn, error (default: info)

View file

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

View file

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

View file

@ -60,7 +60,7 @@ func TestPostSubnet_Created(t *testing.T) {
Name: "sn-new", Name: "sn-new",
VPC: "vpc-1", VPC: "vpc-1",
IfaceType: "vms", IfaceType: "vms",
GatewayIP: "10.0.0.1", InterfaceIP: "10.0.0.1",
CIDR: "10.0.0.0/24", CIDR: "10.0.0.0/24",
} }
body, _ := json.Marshal(req) body, _ := json.Marshal(req)
@ -95,7 +95,7 @@ func TestPostSubnet_IfaceTypeOptional(t *testing.T) {
req := SubnetCreateRequest{ req := SubnetCreateRequest{
Name: "sn-opt", Name: "sn-opt",
VPC: "vpc-1", VPC: "vpc-1",
GatewayIP: "10.0.0.1", InterfaceIP: "10.0.0.1",
CIDR: "10.0.0.0/24", CIDR: "10.0.0.0/24",
// IfaceType omis — doit utiliser default_interface // IfaceType omis — doit utiliser default_interface
} }
@ -113,7 +113,7 @@ func TestPostSubnet_VPCNotFound(t *testing.T) {
Name: "sn-1", Name: "sn-1",
VPC: "vpc-inexistant", VPC: "vpc-inexistant",
IfaceType: "vms", IfaceType: "vms",
GatewayIP: "10.0.0.1", InterfaceIP: "10.0.0.1",
CIDR: "10.0.0.0/24", CIDR: "10.0.0.0/24",
} }
body, _ := json.Marshal(req) body, _ := json.Marshal(req)
@ -132,7 +132,7 @@ func TestPostSubnet_Duplicate(t *testing.T) {
Name: "sn-exist", Name: "sn-exist",
VPC: "vpc-1", VPC: "vpc-1",
IfaceType: "vms", IfaceType: "vms",
GatewayIP: "10.0.0.1", InterfaceIP: "10.0.0.1",
CIDR: "10.0.0.0/24", CIDR: "10.0.0.0/24",
} }
body, _ := json.Marshal(req) body, _ := json.Marshal(req)
@ -150,7 +150,52 @@ func TestPostSubnet_VPCDeleting(t *testing.T) {
Name: "sn-1", Name: "sn-1",
VPC: "vpc-dying", VPC: "vpc-dying",
IfaceType: "vms", 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", CIDR: "10.0.0.0/24",
} }
body, _ := json.Marshal(req) 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/state", "created")
kv.AddInDB(db, "subnet/sn-1/vpc", "vpc-1") 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/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) req := httptest.NewRequest(http.MethodGet, "/subnets/sn-1", nil)
w := httptest.NewRecorder() w := httptest.NewRecorder()
s.SubnetByNameHandler(w, req) s.SubnetByNameHandler(w, req)

View file

@ -44,14 +44,18 @@ func (s *Server) listSubnets(w http.ResponseWriter, _ *http.Request) {
subnets[name].State = value subnets[name].State = value
case "vpc": case "vpc":
subnets[name].VPC = value subnets[name].VPC = value
case "mode":
subnets[name].Mode = value
case "vxlan_id": case "vxlan_id":
subnets[name].VxlanID, _ = strconv.Atoi(value) subnets[name].VxlanID, _ = strconv.Atoi(value)
case "local_iface": case "local_iface":
subnets[name].LocalIface = value subnets[name].LocalIface = value
case "gateway_ip": case "interface_ip":
subnets[name].GatewayIP = value subnets[name].InterfaceIP = value
case "cidr": case "cidr":
subnets[name].CIDR = value subnets[name].CIDR = value
case "default_route":
subnets[name].DefaultRoute = value == "true"
} }
} }
result := make([]Subnet, 0, len(subnets)) 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"}) json.NewEncoder(w).Encode(ErrorResponse{Error: "invalid request body"})
return return
} }
if req.Name == "" || req.VPC == "" || req.GatewayIP == "" || req.CIDR == "" { if req.Name == "" || req.VPC == "" || req.InterfaceIP == "" || req.CIDR == "" {
w.WriteHeader(http.StatusBadRequest) 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 return
} }
cmd := dispatcher.CreateSubnetCommand{ cmd := dispatcher.CreateSubnetCommand{
Name: req.Name, Name: req.Name,
VPC: req.VPC, VPC: req.VPC,
VxlanID: req.VxlanID, Mode: req.Mode,
IfaceType: req.IfaceType, VxlanID: req.VxlanID,
GatewayIP: req.GatewayIP, IfaceType: req.IfaceType,
CIDR: req.CIDR, InterfaceIP: req.InterfaceIP,
CIDR: req.CIDR,
DefaultRoute: req.DefaultRoute,
} }
if err := s.dispatcher.Prepare(cmd); err != nil { if err := s.dispatcher.Prepare(cmd); err != nil {
if _, dbErr := kv.GetFromDB(s.db, "subnet/"+req.Name+"/state"); dbErr == 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 sub.State = value
case "vpc": case "vpc":
sub.VPC = value sub.VPC = value
case "mode":
sub.Mode = value
case "vxlan_id": case "vxlan_id":
sub.VxlanID, _ = strconv.Atoi(value) sub.VxlanID, _ = strconv.Atoi(value)
case "local_iface": case "local_iface":
sub.LocalIface = value sub.LocalIface = value
case "gateway_ip": case "interface_ip":
sub.GatewayIP = value sub.InterfaceIP = value
case "cidr": case "cidr":
sub.CIDR = value sub.CIDR = value
case "default_route":
sub.DefaultRoute = value == "true"
} }
} }
w.WriteHeader(http.StatusAccepted) 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.MetadataPort = entries[prefix+"metadata_port"]
vm.Memory, _ = strconv.Atoi(entries[prefix+"memory"]) vm.Memory, _ = strconv.Atoi(entries[prefix+"memory"])
vm.CPUs, _ = strconv.Atoi(entries[prefix+"cpus"]) vm.CPUs, _ = strconv.Atoi(entries[prefix+"cpus"])
vm.UEFI = entries[prefix+"uefi"] == "true"
subnet := entries[prefix+"subnet"] subnet := entries[prefix+"subnet"]
ip := entries[prefix+"ip"] 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}} vm.Interfaces = []VMInterface{{Subnet: subnet, IP: ip, Primary: true}}
} }
if path := entries[prefix+"volume_path"]; path != "" { diskPrefix := prefix + "disk/"
vm.Storage = []VMStorage{{Path: path}} 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 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"}) json.NewEncoder(w).Encode(ErrorResponse{Error: "invalid request body"})
return 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) 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 return
} }
@ -77,16 +77,21 @@ func (s *Server) startVM(w http.ResponseWriter, r *http.Request) {
return 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{ cmd := dispatcher.StartVMCommand{
Name: req.Name, Name: req.Name,
Subnet: primary.Subnet, Subnet: primary.Subnet,
IP: primary.IP, IP: primary.IP,
MetadataPort: req.MetadataPort, Disks: disks,
VolumePath: req.Storage[0].Path, Memory: req.Memory,
Memory: req.Memory, CPUs: req.CPUs,
CPUs: req.CPUs, UEFI: req.UEFI,
Password: req.Password, Password: req.Password,
SSHKey: req.SSHKey, SSHKey: req.SSHKey,
} }
if err := s.dispatcher.Prepare(cmd); err != nil { 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"}) json.NewEncoder(w).Encode(ErrorResponse{Error: "vpc not found"})
return return
} }
cidr, _ := kv.GetFromDB(s.db, "vpc/"+name+"/cidr")
w.WriteHeader(http.StatusOK) 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) { 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) { func TestPostVpc_Created(t *testing.T) {
s, _ := newTestServer(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() w := httptest.NewRecorder()
s.VpcsHandler(w, httptest.NewRequest(http.MethodPost, "/vpcs", bytes.NewReader(body))) s.VpcsHandler(w, httptest.NewRequest(http.MethodPost, "/vpcs", bytes.NewReader(body)))
if w.Code != http.StatusAccepted { if w.Code != http.StatusAccepted {
@ -67,11 +67,34 @@ func TestPostVpc_Created(t *testing.T) {
if result.State != "creating" { if result.State != "creating" {
t.Errorf("state attendu creating, obtenu %q", result.State) 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) { func TestPostVpc_MissingName(t *testing.T) {
s, _ := newTestServer(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() w := httptest.NewRecorder()
s.VpcsHandler(w, httptest.NewRequest(http.MethodPost, "/vpcs", bytes.NewReader(body))) s.VpcsHandler(w, httptest.NewRequest(http.MethodPost, "/vpcs", bytes.NewReader(body)))
if w.Code != http.StatusBadRequest { if w.Code != http.StatusBadRequest {
@ -82,7 +105,7 @@ func TestPostVpc_MissingName(t *testing.T) {
func TestPostVpc_Duplicate(t *testing.T) { func TestPostVpc_Duplicate(t *testing.T) {
s, db := newTestServer(t) s, db := newTestServer(t)
kv.AddInDB(db, "vpc/vpc-exist/state", "created") 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() w := httptest.NewRecorder()
s.VpcsHandler(w, httptest.NewRequest(http.MethodPost, "/vpcs", bytes.NewReader(body))) s.VpcsHandler(w, httptest.NewRequest(http.MethodPost, "/vpcs", bytes.NewReader(body)))
if w.Code != http.StatusConflict { if w.Code != http.StatusConflict {

View file

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

@ -28,6 +28,22 @@ type Config struct {
Level string `mapstructure:"level"` Level string `mapstructure:"level"`
Debug bool `mapstructure:"debug"` Debug bool `mapstructure:"debug"`
} `mapstructure:"logger"` } `mapstructure:"logger"`
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"` DefaultInterface string `mapstructure:"default_interface"`
Interfaces map[string]string `mapstructure:"interfaces"` Interfaces map[string]string `mapstructure:"interfaces"`
} }
@ -46,6 +62,16 @@ func LoadConfig(path string) (*Config, error) {
v.SetDefault("worker.buffer_size", 100) v.SetDefault("worker.buffer_size", 100)
v.SetDefault("dispatcher.timeout_seconds", 300) v.SetDefault("dispatcher.timeout_seconds", 300)
v.SetDefault("dispatcher.poll_seconds", 2) 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("default_interface", "br-000000")
v.SetDefault("logger.level", "info") v.SetDefault("logger.level", "info")
v.SetDefault("logger.debug", false) 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 { func newConf(t *testing.T, cidr string) Config {
t.Helper() t.Helper()
_, network, _ := net.ParseCIDR(cidr) _, network, _ := net.ParseCIDR(cidr)
_, vpcNet, _ := net.ParseCIDR("10.0.0.0/16")
gw := net.ParseIP("192.168.1.1").To4()
return Config{ return Config{
Network: network, Network: network,
Gateway: net.ParseIP("192.168.1.1").To4(), VPCGateway: gw,
Name: "test", VPCRoute: vpcNet,
ConfDir: t.TempDir(), 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") conf := newConf(t, "192.168.1.0/29")
path, _, _ := GenerateConfig(conf) path, _, _ := GenerateConfig(conf)
content, _ := os.ReadFile(path) content, _ := os.ReadFile(path)
if !strings.Contains(string(content), "dhcp-option=3,192.168.1.1") { 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") _, network, _ := net.ParseCIDR("10.10.0.0/24")
conf := Config{ conf := Config{
Network: network, Network: network,
Gateway: net.ParseIP("10.10.0.1").To4(),
Name: "vpc1", Name: "vpc1",
ConfDir: t.TempDir(), ConfDir: t.TempDir(),
} }
@ -144,7 +179,6 @@ func TestGenerateConfig_CreatesConfDir(t *testing.T) {
_, network, _ := net.ParseCIDR("10.0.0.0/30") _, network, _ := net.ParseCIDR("10.0.0.0/30")
conf := Config{ conf := Config{
Network: network, Network: network,
Gateway: net.ParseIP("10.0.0.1").To4(),
Name: "net", Name: "net",
ConfDir: dir, ConfDir: dir,
} }

View file

@ -14,7 +14,12 @@ func GenerateConfig(c Config) (string, map[string]string, error) {
var sb strings.Builder var sb strings.Builder
fmt.Fprintf(&sb, "no-resolv\n") 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-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") fmt.Fprintf(&sb, "dhcp-option=6,1.1.1.1,8.8.8.8\n\n")
entries := make(map[string]string) entries := make(map[string]string)

View file

@ -5,8 +5,10 @@ import (
) )
type Config struct { type Config struct {
Network *net.IPNet Network *net.IPNet
Gateway net.IP VPCGateway net.IP // next-hop for VPCRoute (option 121)
Name string VPCRoute *net.IPNet // if non-nil, emit dhcp-option=121,VPCRoute,VPCGateway
ConfDir string 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 { type CreateSubnetCommand struct {
Name string Name string
VPC string VPC string
VxlanID int Mode string
IfaceType string VxlanID int
GatewayIP string IfaceType string
CIDR string InterfaceIP string
CIDR string
DefaultRoute bool
} }
func (c CreateSubnetCommand) Prepare(db *badger.DB, cfg *configuration.Config) error { 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 { if _, err := kv.GetFromDB(db, "subnet/"+c.Name+"/state"); err == nil {
return fmt.Errorf("subnet %q already exists", c.Name) 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+"/state", "creating")
kv.AddInDB(db, "subnet/"+c.Name+"/vpc", c.VPC) 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+"/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+"/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 return nil
} }

View file

@ -20,7 +20,7 @@ func TestCreateSubnetCommand_Prepare_Success(t *testing.T) {
kv.AddInDB(db, "vpc/vpc-1/state", "created") kv.AddInDB(db, "vpc/vpc-1/state", "created")
cmd := CreateSubnetCommand{ cmd := CreateSubnetCommand{
Name: "sn-1", VPC: "vpc-1", VxlanID: 100, 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 { if err := cmd.Prepare(db, testCfg()); err != nil {
t.Fatalf("Prepare a échoué : %v", err) 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") kv.AddInDB(db, "vpc/vpc-1/state", "created")
cmd := CreateSubnetCommand{ cmd := CreateSubnetCommand{
Name: "sn-1", VPC: "vpc-1", VxlanID: 100, 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()) cmd.Prepare(db, testCfg())
iface, _ := kv.GetFromDB(db, "subnet/sn-1/local_iface") 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") kv.AddInDB(db, "vpc/vpc-1/state", "created")
cmd := CreateSubnetCommand{ cmd := CreateSubnetCommand{
Name: "sn-1", VPC: "vpc-1", VxlanID: 100, 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()) cmd.Prepare(db, testCfg())
iface, _ := kv.GetFromDB(db, "subnet/sn-1/local_iface") 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") kv.AddInDB(db, "subnet/sn-exist/state", "created")
cmd := CreateSubnetCommand{ cmd := CreateSubnetCommand{
Name: "sn-exist", VPC: "vpc-1", VxlanID: 100, 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 { if err := cmd.Prepare(db, testCfg()); err == nil {
t.Error("Prepare devrait échouer sur un subnet déjà existant") 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) _, db := newTestDispatcher(t)
cmd := CreateSubnetCommand{ cmd := CreateSubnetCommand{
Name: "sn-1", VPC: "vpc-inexistant", VxlanID: 100, 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 { if err := cmd.Prepare(db, testCfg()); err == nil {
t.Error("Prepare devrait échouer si le VPC n'existe pas") 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") kv.AddInDB(db, "vpc/vpc-dying/state", "deleting")
cmd := CreateSubnetCommand{ cmd := CreateSubnetCommand{
Name: "sn-1", VPC: "vpc-dying", VxlanID: 100, 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 { if err := cmd.Prepare(db, testCfg()); err == nil {
t.Error("Prepare devrait échouer si le VPC est en cours de suppression") 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") kv.AddInDB(db, "vpc/vpc-gone/state", "deleted")
cmd := CreateSubnetCommand{ cmd := CreateSubnetCommand{
Name: "sn-1", VPC: "vpc-gone", VxlanID: 100, 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 { if err := cmd.Prepare(db, testCfg()); err == nil {
t.Error("Prepare devrait échouer si le VPC est supprimé") 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 --- // --- DeleteSubnetCommand.Prepare ---
func TestDeleteSubnetCommand_Prepare_Success(t *testing.T) { func TestDeleteSubnetCommand_Prepare_Success(t *testing.T) {

View file

@ -2,7 +2,9 @@ package dispatcher
import ( import (
"fmt" "fmt"
"math/rand"
"strconv" "strconv"
"strings"
"time" "time"
configuration "git.g3e.fr/syonad/two/internal/config/agent" configuration "git.g3e.fr/syonad/two/internal/config/agent"
@ -11,16 +13,21 @@ import (
"github.com/dgraph-io/badger/v4" "github.com/dgraph-io/badger/v4"
) )
type VMDisk struct {
Path string
Dev string
}
type StartVMCommand struct { type StartVMCommand struct {
Name string Name string
Subnet string Subnet string
IP string IP string
MetadataPort string Disks []VMDisk
VolumePath string Memory int
Memory int CPUs int
CPUs int UEFI bool
Password string Password string
SSHKey string SSHKey string
} }
func (c StartVMCommand) Prepare(db *badger.DB, _ *configuration.Config) error { 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" { if subnetState == "deleting" || subnetState == "deleted" {
return fmt.Errorf("subnet %q is %s", c.Subnet, subnetState) 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+"/state", "starting")
kv.AddInDB(db, "vm/"+c.Name+"/subnet", c.Subnet) kv.AddInDB(db, "vm/"+c.Name+"/subnet", c.Subnet)
kv.AddInDB(db, "vm/"+c.Name+"/ip", c.IP) 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+"/metadata_port", strconv.Itoa(port))
kv.AddInDB(db, "vm/"+c.Name+"/volume_path", c.VolumePath) 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+"/memory", strconv.Itoa(c.Memory))
kv.AddInDB(db, "vm/"+c.Name+"/cpus", strconv.Itoa(c.CPUs)) 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 != "" { if c.Password != "" {
kv.AddInDB(db, "vm/"+c.Name+"/password", 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 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 { 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 { for {
@ -66,7 +104,7 @@ func (c StartVMCommand) Execute(db *badger.DB, cfg *configuration.Config) error
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) return vm.StartVM(db, c.Name, cfg)
} }
type StopVMCommand struct { type StopVMCommand struct {

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 ( import (
"fmt" "fmt"
"net"
"strings" "strings"
"time" "time"
@ -13,12 +14,19 @@ import (
type CreateVPCCommand struct { type CreateVPCCommand struct {
Name string Name string
CIDR string
} }
func (c CreateVPCCommand) Prepare(db *badger.DB, _ *configuration.Config) error { func (c CreateVPCCommand) Prepare(db *badger.DB, _ *configuration.Config) error {
if _, err := kv.GetFromDB(db, "vpc/"+c.Name+"/state"); err == nil { if _, err := kv.GetFromDB(db, "vpc/"+c.Name+"/state"); err == nil {
return fmt.Errorf("vpc %q already exists", c.Name) 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") return kv.AddInDB(db, "vpc/"+c.Name+"/state", "creating")
} }

View file

@ -10,7 +10,7 @@ import (
func TestCreateVPCCommand_Prepare_NewVPC(t *testing.T) { func TestCreateVPCCommand_Prepare_NewVPC(t *testing.T) {
_, db := newTestDispatcher(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 { if err := cmd.Prepare(db, nil); err != nil {
t.Fatalf("Prepare a échoué : %v", err) t.Fatalf("Prepare a échoué : %v", err)
} }
@ -21,17 +21,32 @@ func TestCreateVPCCommand_Prepare_NewVPC(t *testing.T) {
if state != "creating" { if state != "creating" {
t.Errorf("state attendu creating, obtenu %q", state) 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) { func TestCreateVPCCommand_Prepare_Duplicate(t *testing.T) {
_, db := newTestDispatcher(t) _, db := newTestDispatcher(t)
kv.AddInDB(db, "vpc/vpc-exist/state", "created") 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 { if err := cmd.Prepare(db, nil); err == nil {
t.Error("Prepare devrait échouer sur un VPC déjà existant") 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 --- // --- DeleteVPCCommand.Prepare ---
func TestDeleteVPCCommand_Prepare_Success(t *testing.T) { 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() 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", if err := addRule("FORWARD",
"--out-interface", bridge, "--out-interface", iface,
"-p", "arp", "-p", "arp",
"--arp-op", "Request", "--arp-op", "Request",
"--arp-ip-dst", gatewayIP, "--arp-ip-dst", ip,
"-j", "DROP"); err != nil { "-j", "DROP"); err != nil {
return fmt.Errorf("ebtables arp rule: %w", err) return fmt.Errorf("ebtables arp rule: %w", err)
} }
return nil return nil
} }
func DropDHCP(bridge string) error { func DropDHCP(iface, ip string) error {
if err := addRule("FORWARD", if err := addRule("FORWARD",
"--out-interface", bridge, "--out-interface", iface,
"-p", "IPv4", "-p", "IPv4",
"--ip-protocol", "udp", "--ip-protocol", "udp",
"--ip-source-port", "67:68", "--ip-source-port", "67:68",
"--ip-destination-port", "67:68", "--ip-destination-port", "67:68",
"--ip-source", ip,
"-j", "DROP"); err != nil { "-j", "DROP"); err != nil {
return fmt.Errorf("ebtables dhcp rule: %w", err) return fmt.Errorf("ebtables dhcp rule: %w", err)
} }
return nil return nil
} }
func DeleteARPToGateway(bridge, gatewayIP string) error { func DeleteARPToGateway(iface, ip string) error {
return deleteRule("FORWARD", return deleteRule("FORWARD",
"--out-interface", bridge, "--out-interface", iface,
"-p", "arp", "-p", "arp",
"--arp-op", "Request", "--arp-op", "Request",
"--arp-ip-dst", gatewayIP, "--arp-ip-dst", ip,
"-j", "DROP") "-j", "DROP")
} }
func DeleteDHCP(bridge string) error { func DeleteDHCP(iface, ip string) error {
return deleteRule("FORWARD", return deleteRule("FORWARD",
"--out-interface", bridge, "--out-interface", iface,
"-p", "IPv4", "-p", "IPv4",
"--ip-protocol", "udp", "--ip-protocol", "udp",
"--ip-source-port", "67:68", "--ip-source-port", "67:68",
"--ip-destination-port", "67:68", "--ip-destination-port", "67:68",
"--ip-source", ip,
"-j", "DROP") "-j", "DROP")
} }

View file

@ -3,18 +3,18 @@ package metadata
import ( import (
"fmt" "fmt"
configuration "git.g3e.fr/syonad/two/internal/config/agent"
"git.g3e.fr/syonad/two/pkg/systemd" "git.g3e.fr/syonad/two/pkg/systemd"
"github.com/dgraph-io/badger/v4"
) )
func StartMetadata(config NoCloudConfig, db *badger.DB, dryrun bool) error { func StartMetadata(config NoCloudConfig, cfg *configuration.Config, dryrun bool) error {
service, err := systemd.New() service, err := systemd.New()
if err != nil { if err != nil {
return fmt.Errorf("failed to connect to systemd: %w", err) return fmt.Errorf("failed to connect to systemd: %w", err)
} }
defer service.Close() defer service.Close()
LoadNcCloudInDB(config, db) LoadNcCloudInDB(config, cfg.Metadata.RunDir)
if !dryrun { if !dryrun {
if err := service.Start("metadata@" + config.Name + ".service"); err != nil { if err := service.Start("metadata@" + config.Name + ".service"); err != nil {
return fmt.Errorf("failed to start metadata@%s: %w", config.Name, err) return fmt.Errorf("failed to start metadata@%s: %w", config.Name, err)
@ -23,17 +23,17 @@ func StartMetadata(config NoCloudConfig, db *badger.DB, dryrun bool) error {
return nil return nil
} }
func StopMetadata(vm_name string, db *badger.DB, dryrun bool) error { func StopMetadata(vmName string, cfg *configuration.Config, dryrun bool) error {
service, err := systemd.New() service, err := systemd.New()
if err != nil { if err != nil {
return fmt.Errorf("failed to connect to systemd: %w", err) return fmt.Errorf("failed to connect to systemd: %w", err)
} }
defer service.Close() defer service.Close()
UnLoadNoCloudInDB(vm_name, db) UnLoadNoCloudInDB(vmName, cfg.Metadata.RunDir)
if !dryrun { if !dryrun {
if err := service.Stop("metadata@" + vm_name + ".service"); err != nil { if err := service.Stop("metadata@" + vmName + ".service"); err != nil {
return fmt.Errorf("failed to stop metadata@%s: %w", vm_name, err) return fmt.Errorf("failed to stop metadata@%s: %w", vmName, err)
} }
} }
return nil return nil

View file

@ -3,10 +3,10 @@ package metadata
import ( import (
"net/http" "net/http"
"net/http/httptest" "net/http/httptest"
"os"
"path/filepath"
"strings" "strings"
"testing" "testing"
"git.g3e.fr/syonad/two/pkg/db/kv"
) )
func newCfg() NoCloudConfig { func newCfg() NoCloudConfig {
@ -20,11 +20,9 @@ func newCfg() NoCloudConfig {
} }
} }
func newTestDB(t *testing.T) interface{ Close() error } { func useTestDir(t *testing.T) string {
t.Helper() t.Helper()
db := kv.InitDB(kv.Config{Path: t.TempDir()}, false) return t.TempDir()
t.Cleanup(func() { db.Close() })
return db
} }
// --- RenderConfig --- // --- RenderConfig ---
@ -108,78 +106,67 @@ func TestRenderConfig_SpecialCharsInName(t *testing.T) {
// --- LoadNcCloudInDB / UnLoadNoCloudInDB --- // --- LoadNcCloudInDB / UnLoadNoCloudInDB ---
func TestLoadNcCloudInDB_StoresAllKeys(t *testing.T) { func readTestFile(t *testing.T, dir, vmName, name string) string {
db := kv.InitDB(kv.Config{Path: t.TempDir()}, false) t.Helper()
t.Cleanup(func() { db.Close() }) b, err := os.ReadFile(filepath.Join(dir, vmName, name))
if err != nil {
cfg := newCfg() t.Errorf("fichier %q absent après LoadNcCloudInDB : %v", name, err)
LoadNcCloudInDB(cfg, db) return ""
keys := []string{
"metadata/vm1/meta-data",
"metadata/vm1/user-data",
"metadata/vm1/network-config",
"metadata/vm1/vendor-data",
"metadata/vm1/vpc",
"metadata/vm1/bind_ip",
"metadata/vm1/bind_port",
} }
for _, key := range keys { return string(b)
val, err := kv.GetFromDB(db, key) }
if err != nil {
t.Errorf("clé %q absente après LoadNcCloudInDB : %v", key, err) func TestLoadNcCloudInDB_StoresAllFiles(t *testing.T) {
} dir := useTestDir(t)
if val == "" && key != "metadata/vm1/user-data" { LoadNcCloudInDB(newCfg(), dir)
t.Errorf("clé %q vide après LoadNcCloudInDB", key)
files := []string{"meta-data", "user-data", "network-config", "vendor-data", "vpc", "bind_ip", "bind_port"}
for _, f := range files {
path := filepath.Join(dir, "vm1", f)
if _, err := os.Stat(path); err != nil {
t.Errorf("fichier %q absent : %v", f, err)
} }
} }
} }
func TestLoadNcCloudInDB_VpcAndBindValues(t *testing.T) { func TestLoadNcCloudInDB_VpcAndBindValues(t *testing.T) {
db := kv.InitDB(kv.Config{Path: t.TempDir()}, false) dir := useTestDir(t)
t.Cleanup(func() { db.Close() }) LoadNcCloudInDB(newCfg(), dir)
cfg := newCfg() if vpc := readTestFile(t, dir, "vm1", "vpc"); vpc != "vpc-test" {
LoadNcCloudInDB(cfg, db)
vpc, _ := kv.GetFromDB(db, "metadata/vm1/vpc")
if vpc != "vpc-test" {
t.Errorf("vpc attendu %q, obtenu %q", "vpc-test", vpc) t.Errorf("vpc attendu %q, obtenu %q", "vpc-test", vpc)
} }
if ip := readTestFile(t, dir, "vm1", "bind_ip"); ip != "169.254.169.254" {
ip, _ := kv.GetFromDB(db, "metadata/vm1/bind_ip")
if ip != "169.254.169.254" {
t.Errorf("bind_ip attendu %q, obtenu %q", "169.254.169.254", ip) t.Errorf("bind_ip attendu %q, obtenu %q", "169.254.169.254", ip)
} }
if port := readTestFile(t, dir, "vm1", "bind_port"); port != "80" {
port, _ := kv.GetFromDB(db, "metadata/vm1/bind_port")
if port != "80" {
t.Errorf("bind_port attendu %q, obtenu %q", "80", port) t.Errorf("bind_port attendu %q, obtenu %q", "80", port)
} }
} }
func TestUnLoadNoCloudInDB_RemovesAllKeys(t *testing.T) { func TestUnLoadNoCloudInDB_RemovesAllFiles(t *testing.T) {
db := kv.InitDB(kv.Config{Path: t.TempDir()}, false) dir := useTestDir(t)
t.Cleanup(func() { db.Close() }) LoadNcCloudInDB(newCfg(), dir)
UnLoadNoCloudInDB("vm1", dir)
cfg := newCfg() if _, err := os.Stat(filepath.Join(dir, "vm1")); !os.IsNotExist(err) {
LoadNcCloudInDB(cfg, db) t.Error("répertoire vm1 devrait être supprimé après UnLoadNoCloudInDB")
UnLoadNoCloudInDB("vm1", db)
keys := []string{
"metadata/vm1/meta-data",
"metadata/vm1/user-data",
"metadata/vm1/network-config",
"metadata/vm1/vendor-data",
"metadata/vm1/vpc",
"metadata/vm1/bind_ip",
"metadata/vm1/bind_port",
} }
for _, key := range keys { }
_, err := kv.GetFromDB(db, key)
if err == nil { func TestUnLoadNoCloudInDB_DoesNotAffectOtherVMs(t *testing.T) {
t.Errorf("clé %q devrait être supprimée après UnLoadNoCloudInDB", key) dir := useTestDir(t)
}
cfg1 := newCfg()
cfg2 := newCfg()
cfg2.Name = "vm2"
LoadNcCloudInDB(cfg1, dir)
LoadNcCloudInDB(cfg2, dir)
UnLoadNoCloudInDB("vm1", dir)
if _, err := os.Stat(filepath.Join(dir, "vm2", "vpc")); err != nil {
t.Errorf("vm2 ne devrait pas être supprimée : %v", err)
} }
} }
@ -277,23 +264,3 @@ func TestRootHandler_ContentType(t *testing.T) {
t.Errorf("Content-Type attendu text/yaml, obtenu %q", ct) t.Errorf("Content-Type attendu text/yaml, obtenu %q", ct)
} }
} }
// --- UnLoadNoCloudInDB_DoesNotAffectOtherVMs ---
func TestUnLoadNoCloudInDB_DoesNotAffectOtherVMs(t *testing.T) {
db := kv.InitDB(kv.Config{Path: t.TempDir()}, false)
t.Cleanup(func() { db.Close() })
cfg1 := newCfg()
cfg2 := newCfg()
cfg2.Name = "vm2"
LoadNcCloudInDB(cfg1, db)
LoadNcCloudInDB(cfg2, db)
UnLoadNoCloudInDB("vm1", db)
_, err := kv.GetFromDB(db, "metadata/vm2/vpc")
if err != nil {
t.Errorf("vm2 ne devrait pas être supprimée : %v", err)
}
}

View file

@ -3,10 +3,9 @@ package metadata
import ( import (
"bytes" "bytes"
"embed" "embed"
"os"
"path/filepath"
"text/template" "text/template"
"git.g3e.fr/syonad/two/pkg/db/kv"
"github.com/dgraph-io/badger/v4"
) )
//go:embed templates/*.tmpl //go:embed templates/*.tmpl
@ -26,21 +25,25 @@ func RenderConfig(path string, cfg NoCloudConfig) (string, error) {
return buf.String(), nil return buf.String(), nil
} }
func LoadNcCloudInDB(config NoCloudConfig, db *badger.DB) { func LoadNcCloudInDB(config NoCloudConfig, runDir string) {
meta_data, _ := RenderConfig("templates/meta-data.tmpl", config) meta_data, _ := RenderConfig("templates/meta-data.tmpl", config)
user_data, _ := RenderConfig("templates/user-data.tmpl", config) user_data, _ := RenderConfig("templates/user-data.tmpl", config)
network_config, _ := RenderConfig("templates/network-config.tmpl", config) network_config, _ := RenderConfig("templates/network-config.tmpl", config)
vendor_data, _ := RenderConfig("templates/vendor-data.tmpl", config) vendor_data, _ := RenderConfig("templates/vendor-data.tmpl", config)
kv.AddInDB(db, "metadata/"+config.Name+"/meta-data", meta_data) dir := filepath.Join(runDir, config.Name)
kv.AddInDB(db, "metadata/"+config.Name+"/user-data", user_data) if err := os.MkdirAll(dir, 0755); err != nil {
kv.AddInDB(db, "metadata/"+config.Name+"/network-config", network_config) return
kv.AddInDB(db, "metadata/"+config.Name+"/vendor-data", vendor_data) }
kv.AddInDB(db, "metadata/"+config.Name+"/vpc", config.VpcName) os.WriteFile(filepath.Join(dir, "meta-data"), []byte(meta_data), 0644)
kv.AddInDB(db, "metadata/"+config.Name+"/bind_ip", config.BindIP) os.WriteFile(filepath.Join(dir, "user-data"), []byte(user_data), 0644)
kv.AddInDB(db, "metadata/"+config.Name+"/bind_port", config.BindPort) os.WriteFile(filepath.Join(dir, "network-config"), []byte(network_config), 0644)
os.WriteFile(filepath.Join(dir, "vendor-data"), []byte(vendor_data), 0644)
os.WriteFile(filepath.Join(dir, "vpc"), []byte(config.VpcName), 0644)
os.WriteFile(filepath.Join(dir, "bind_ip"), []byte(config.BindIP), 0644)
os.WriteFile(filepath.Join(dir, "bind_port"), []byte(config.BindPort), 0644)
} }
func UnLoadNoCloudInDB(vm_name string, db *badger.DB) { func UnLoadNoCloudInDB(vmName string, runDir string) {
kv.DeleteInDB(db, "metadata/"+vm_name) os.RemoveAll(filepath.Join(runDir, vmName))
} }

View file

@ -5,12 +5,13 @@ import (
"log" "log"
"net" "net"
"net/http" "net/http"
"os"
"path/filepath"
"strconv" "strconv"
"strings"
"time" "time"
configuration "git.g3e.fr/syonad/two/internal/config/agent"
"git.g3e.fr/syonad/two/internal/netns" "git.g3e.fr/syonad/two/internal/netns"
"git.g3e.fr/syonad/two/pkg/db/kv"
) )
var data NoCloudData var data NoCloudData
@ -23,47 +24,23 @@ func getIP(r *http.Request) string {
return ip return ip
} }
func getFromDB(config ServerConfig) NoCloudData { func readFile(dir, name string) string {
var netns_name string b, _ := os.ReadFile(filepath.Join(dir, name))
var port int return strings.TrimRight(string(b), "\n")
var iface string }
conf_db, _ := configuration.LoadConfig(config.ConfFile) func getFromFiles(config ServerConfig) NoCloudData {
dir := filepath.Join(config.RunDir, config.VmName)
db := kv.InitDB(kv.Config{Path: conf_db.Database.Path}, true) port, _ := strconv.Atoi(readFile(dir, "bind_port"))
defer db.Close()
metadata, _ := kv.GetFromDB(db, "metadata/"+config.VmName+"/meta-data")
userdata, _ := kv.GetFromDB(db, "metadata/"+config.VmName+"/user-data")
networkconfig, _ := kv.GetFromDB(db, "metadata/"+config.VmName+"/network-config")
vendordata, _ := kv.GetFromDB(db, "metadata/"+config.VmName+"/vendor-data")
if config.Netns == "" {
netns_name, _ = kv.GetFromDB(db, "metadata/"+config.VmName+"/vpc")
} else {
netns_name = config.Netns
}
if config.Iface == "" {
iface, _ = kv.GetFromDB(db, "metadata/"+config.VmName+"/bind_ip")
} else {
iface = config.Iface
}
if config.Port == 0 {
sport, _ := kv.GetFromDB(db, "metadata/"+config.VmName+"/bind_port")
port, _ = strconv.Atoi(sport)
} else {
port = config.Port
}
return NoCloudData{ return NoCloudData{
MetaData: metadata, MetaData: readFile(dir, "meta-data"),
UserData: userdata, UserData: readFile(dir, "user-data"),
NetworkConfig: networkconfig, NetworkConfig: readFile(dir, "network-config"),
VendorData: vendordata, VendorData: readFile(dir, "vendor-data"),
NetNs: netns_name, NetNs: readFile(dir, "vpc"),
Iface: iface, Iface: readFile(dir, "bind_ip"),
Port: port, Port: port,
} }
} }
@ -93,7 +70,7 @@ func rootHandler(w http.ResponseWriter, r *http.Request) {
} }
func StartServer(config ServerConfig) { func StartServer(config ServerConfig) {
data = getFromDB(config) data = getFromFiles(config)
if data.NetNs != "" { if data.NetNs != "" {
if err := netns.Enter(data.NetNs); err != nil { if err := netns.Enter(data.NetNs); err != nil {

View file

@ -11,12 +11,8 @@ type NoCloudData struct {
} }
type ServerConfig struct { type ServerConfig struct {
Netns string VmName string
File string RunDir string
Iface string
Port int
ConfFile string
VmName string
} }
type NoCloudConfig struct { type NoCloudConfig struct {

View file

@ -2,12 +2,8 @@
users: users:
- name: syonad - name: syonad
lock_passwd: false lock_passwd: false
gecos: alpine Cloud User
groups: [adm, wheel]
doas:
- permit nopass syonad
sudo: ["ALL=(ALL) NOPASSWD:ALL"] sudo: ["ALL=(ALL) NOPASSWD:ALL"]
shell: /bin/ash shell: /bin/bash
passwd: "{{ .Password }}" passwd: "{{ .Password }}"
ssh_authorized_keys: 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 ( import (
"fmt" "fmt"
"os"
"os/exec" "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 { func Start(cfg Config) error {
memory := cfg.Memory memory := cfg.Memory
if memory == 0 { if memory == 0 {
@ -27,20 +22,83 @@ func Start(cfg Config) error {
cpus = 1 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", "-enable-kvm",
"-cpu", "host", "-cpu", "host",
"-m", fmt.Sprintf("%d", memory), "-m", fmt.Sprintf("%d", memory),
"-smp", fmt.Sprintf("%d", cpus), "-smp", fmt.Sprintf("%d", cpus),
"-serial", fmt.Sprintf("unix:/tmp/%s.sock,server,nowait", cfg.Name), "-serial", fmt.Sprintf("unix:%s,server,nowait", serialSock),
"-monitor", fmt.Sprintf("unix:/tmp/%s.mon-sock,server,nowait", cfg.Name), "-monitor", fmt.Sprintf("unix:%s,server,nowait", monitorSock),
"-qmp", fmt.Sprintf("unix:/tmp/%s.qmp-sock,server,nowait", cfg.Name), "-qmp", fmt.Sprintf("unix:%s,server,nowait", qmpSock),
"-display", "none", "-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), "-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), "-device", fmt.Sprintf("virtio-net-pci,netdev=net0,mac=%s", cfg.Mac),
"-daemonize", "-daemonize",
) )
cmd := exec.Command("qemu-system-x86_64", args...)
if err := cmd.Run(); err != nil { if err := cmd.Run(); err != nil {
return fmt.Errorf("qemu-system-x86_64: %w", err) return fmt.Errorf("qemu-system-x86_64: %w", err)
} }

View file

@ -2,12 +2,9 @@
package qemu package qemu
import "errors" import (
"errors"
type Config struct { )
Name, Mac, VolumePath string
TapID, Memory, CPUs int
}
func Start(_ Config) error { func Start(_ Config) error {
return errors.New("vm: not supported on this platform") 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 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) return fmt.Errorf("create veth: %w", err)
} }
if err := netif.CreateBridge(d.bridge, 1500); err != nil { switch d.mode {
return fmt.Errorf("create bridge: %w", err) case "vxlan":
} if err := setupVxlanHost(d, vethE); err != nil {
return 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)
} }
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 { 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 { if err := netif.LinkSetUp(iface); err != nil {
return fmt.Errorf("set up %s: %w", iface, err) 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) return fmt.Errorf("set up interfaces in netns: %w", err)
} }
if err := netns.Call(d.vpc, func() error { switch d.mode {
return netif.AddrAdd(d.bridge, d.gatewayIP) case "vxlan":
}); err != nil { if err := netns.Call(d.vpc, func() error {
return fmt.Errorf("add addr to bridge in netns: %w", err) 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 startDHCP(db, subnetName, d)
return netif.RouteAdd(d.bridge, d.cidr) }
}); err != nil {
return fmt.Errorf("add route in netns: %w", err)
}
if err := ebtables.DropARPToGateway(d.bridge, d.gatewayIP.String()); err != nil { func setupVxlanHost(d subnetData, vethE string) error {
return err vxlanIface := fmt.Sprintf("vxlan-%d", d.vxlanID)
}
if err := ebtables.DropDHCP(d.bridge); err != nil {
return err
}
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{ conf := dhcp.Config{
Network: d.cidr, Network: d.cidr,
Gateway: d.gatewayIP,
Name: d.vpc + "_" + d.bridge, Name: d.vpc + "_" + d.bridge,
ConfDir: "/etc/dnsmasq.d", 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) _, entries, err := dhcp.GenerateConfig(conf)
if err != nil { if err != nil {
return fmt.Errorf("generate dhcp config: %w", err) 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 { if err := svc.Start("dnsmasq@" + conf.Name + ".service"); err != nil {
return fmt.Errorf("start dnsmasq: %w", err) return fmt.Errorf("start dnsmasq: %w", err)
} }
return nil
return kv.AddInDB(db, "subnet/"+subnetName+"/state", "created")
} }

View file

@ -11,13 +11,16 @@ import (
) )
type subnetData struct { type subnetData struct {
vpc string vpc string
subnetID string subnetID string
bridge string bridge string
vxlanID int mode string
localIface string vxlanID int
gatewayIP net.IP localIface string
cidr *net.IPNet interfaceIP net.IP
cidr *net.IPNet
vpcCIDR *net.IPNet
defaultRoute bool
} }
func loadSubnet(db *badger.DB, name string) (subnetData, error) { 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 d.vpc = vpc
vxlanIDStr, err := kv.GetFromDB(db, "subnet/"+name+"/vxlan_id") mode, err := kv.GetFromDB(db, "subnet/"+name+"/mode")
if err != nil { 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) d.mode = mode
if err != nil {
return d, fmt.Errorf("parse vxlan_id: %w", err) 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") localIface, err := kv.GetFromDB(db, "subnet/"+name+"/local_iface")
if err != nil { if err != nil {
@ -48,15 +59,15 @@ func loadSubnet(db *badger.DB, name string) (subnetData, error) {
} }
d.localIface = localIface d.localIface = localIface
gatewayIPStr, err := kv.GetFromDB(db, "subnet/"+name+"/gateway_ip") interfaceIPStr, err := kv.GetFromDB(db, "subnet/"+name+"/interface_ip")
if err != nil { 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) interfaceIP := net.ParseIP(interfaceIPStr)
if gatewayIP == nil { if interfaceIP == nil {
return d, fmt.Errorf("invalid gateway_ip: %s", gatewayIPStr) return d, fmt.Errorf("invalid interface_ip: %s", interfaceIPStr)
} }
d.gatewayIP = gatewayIP d.interfaceIP = interfaceIP
cidrStr, err := kv.GetFromDB(db, "subnet/"+name+"/cidr") cidrStr, err := kv.GetFromDB(db, "subnet/"+name+"/cidr")
if err != nil { if err != nil {
@ -68,5 +79,21 @@ func loadSubnet(db *badger.DB, name string) (subnetData, error) {
} }
d.cidr = ipNet 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 return d, nil
} }

View file

@ -27,36 +27,64 @@ func DeleteSubnet(db *badger.DB, subnetName string) error {
return err 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() svc, err := systemd.New()
if err != nil { if err != nil {
return fmt.Errorf("connect to systemd: %w", err) return fmt.Errorf("connect to systemd: %w", err)
} }
defer svc.Close() defer svc.Close()
if err := svc.Stop("dnsmasq@" + d.vpc + "_" + d.bridge + ".service"); err != nil { svcName := "dnsmasq@" + d.vpc + "_" + d.bridge + ".service"
return fmt.Errorf("stop dnsmasq: %w", err) 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) { 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) return fmt.Errorf("remove dnsmasq config: %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)
} }
return nil
}
if err := ebtables.DeleteARPToGateway(d.bridge, d.gatewayIP.String()); err != nil { func deleteSubnetVxlan(d subnetData) error {
return fmt.Errorf("delete ebtables arp rule: %w", err) vxlanIface := fmt.Sprintf("vxlan-%d", d.vxlanID)
} vethI := "v-" + d.subnetID + "-i"
if err := ebtables.DeleteDHCP(d.bridge); err != nil {
return fmt.Errorf("delete ebtables dhcp rule: %w", err)
}
if err := netns.Call(d.vpc, func() error { 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) return netif.DeleteLink(d.bridge)
}); err != nil { }); 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 { 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 { if err := netif.DeleteLink(d.bridge); err != nil {
return fmt.Errorf("delete bridge: %w", err) return fmt.Errorf("delete bridge: %w", err)
} }
return nil
return kv.AddInDB(db, "subnet/"+subnetName+"/state", "deleted") }
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,7 +2,11 @@ package vm
import ( import (
"fmt" "fmt"
"io"
"os"
"path/filepath"
configuration "git.g3e.fr/syonad/two/internal/config/agent"
"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"
@ -13,7 +17,7 @@ import (
"github.com/dgraph-io/badger/v4" "github.com/dgraph-io/badger/v4"
) )
func StartVM(db *badger.DB, name string) error { func StartVM(db *badger.DB, name string, cfg *configuration.Config) error {
state, err := kv.GetFromDB(db, "vm/"+name+"/state") state, err := kv.GetFromDB(db, "vm/"+name+"/state")
if err != nil { if err != nil {
return err return err
@ -32,7 +36,7 @@ func StartVM(db *badger.DB, name string) error {
} }
if err := netns.Call(d.vpcName, func() 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 { }); err != nil {
return fmt.Errorf("add metadata redirect: %w", err) return fmt.Errorf("add metadata redirect: %w", err)
} }
@ -40,26 +44,67 @@ func StartVM(db *badger.DB, name string) error {
if err := metadata.StartMetadata(metadata.NoCloudConfig{ if err := metadata.StartMetadata(metadata.NoCloudConfig{
Name: name, Name: name,
VpcName: d.vpcName, VpcName: d.vpcName,
BindIP: d.gatewayIP, BindIP: d.interfaceIP,
BindPort: d.metadataPort, BindPort: d.metadataPort,
Password: d.password, Password: d.password,
SSHKEY: d.sshkey, SSHKEY: d.sshkey,
}, db, false); err != nil { }, cfg, false); err != nil {
return fmt.Errorf("start metadata: %w", err) 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 { if err := netns.Call(d.vpcName, func() error {
return qemu.Start(qemu.Config{ return qemu.Start(qcfg)
Name: name,
TapID: d.tapID,
Mac: d.mac,
VolumePath: d.volumePath,
Memory: d.memory,
CPUs: d.cpus,
})
}); err != nil { }); err != nil {
return fmt.Errorf("start qemu: %w", err) return fmt.Errorf("start qemu: %w", err)
} }
return kv.AddInDB(db, "vm/"+name+"/state", "started") 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" "github.com/dgraph-io/badger/v4"
) )
type diskEntry struct {
path string
dev string
}
type vmData struct { type vmData struct {
subnetName string subnetName string
vpcName string vpcName string
gatewayIP string interfaceIP string
bridge string bridge string
tapID int tapID int
ip string ip string
metadataPort string metadataPort string
mac string mac string
volumePath string disks []diskEntry
memory int memory int
cpus int cpus int
uefi bool
password string password string
sshkey string sshkey string
} }
@ -43,11 +49,11 @@ func loadVM(db *badger.DB, name string) (vmData, error) {
} }
d.vpcName = vpcName d.vpcName = vpcName
gatewayIP, err := kv.GetFromDB(db, "subnet/"+subnetName+"/gateway_ip") interfaceIP, err := kv.GetFromDB(db, "subnet/"+subnetName+"/interface_ip")
if err != nil { 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") tapIDStr, err := kv.GetFromDB(db, "vm/"+name+"/tap_id")
if err != nil { if err != nil {
@ -81,11 +87,18 @@ func loadVM(db *badger.DB, name string) (vmData, error) {
} }
d.mac = mac d.mac = mac
volumePath, err := kv.GetFromDB(db, "vm/"+name+"/volume_path") diskEntries, err := kv.ListByPrefix(db, "vm/"+name+"/disk/")
if err != nil { 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") memoryStr, err := kv.GetFromDB(db, "vm/"+name+"/memory")
if err != nil { if err != nil {
@ -105,6 +118,9 @@ func loadVM(db *badger.DB, name string) (vmData, error) {
return d, fmt.Errorf("parse cpus: %w", err) 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.password, _ = kv.GetFromDB(db, "vm/"+name+"/password")
d.sshkey, _ = kv.GetFromDB(db, "vm/"+name+"/sshkey") d.sshkey, _ = kv.GetFromDB(db, "vm/"+name+"/sshkey")

View file

@ -2,6 +2,8 @@ package vm
import ( import (
"fmt" "fmt"
"os"
"path/filepath"
"time" "time"
configuration "git.g3e.fr/syonad/two/internal/config/agent" configuration "git.g3e.fr/syonad/two/internal/config/agent"
@ -29,35 +31,28 @@ func StopVM(db *badger.DB, name string, cfg *configuration.Config) error {
return err 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 { if _, err := os.Stat(socketPath); err == nil {
return fmt.Errorf("qmp system_powerdown: %w", err) // socket présent : tenter l'arrêt gracieux
} if _, err := qmp.Send(socketPath, []string{`{"execute":"system_powerdown"}`}); err == nil {
waitQMPDead(socketPath,
// attendre l'arrêt effectif de la VM ; forcer via quit après timeout time.Duration(cfg.Dispatcher.TimeoutSeconds)*time.Second,
timeout := time.After(time.Duration(cfg.Dispatcher.TimeoutSeconds) * time.Second) time.Duration(cfg.Dispatcher.PollSeconds)*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
}
} }
// 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 { 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 { }); err != nil {
return fmt.Errorf("delete metadata redirect: %w", err) return fmt.Errorf("delete metadata redirect: %w", err)
} }
if err := metadata.StopMetadata(name, db, false); err != nil { if err := metadata.StopMetadata(name, cfg, false); err != nil {
return fmt.Errorf("stop metadata: %w", err) return fmt.Errorf("stop metadata: %w", err)
} }
@ -65,5 +60,25 @@ func StopVM(db *badger.DB, name string, cfg *configuration.Config) error {
return fmt.Errorf("delete tap: %w", err) 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") 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

@ -1,6 +1,8 @@
package kv package kv
import ( import (
"log"
"github.com/dgraph-io/badger/v4" "github.com/dgraph-io/badger/v4"
) )
@ -15,6 +17,7 @@ func InitDB(conf Config, readonly bool) *badger.DB {
opts.NumLevelZeroTablesStall = 2 opts.NumLevelZeroTablesStall = 2
db, err := badger.Open(opts) db, err := badger.Open(opts)
if err != nil { if err != nil {
log.Printf("kv.InitDB (readonly=%v, path=%s): %v", readonly, conf.Path, err)
panic(err) panic(err)
} }
return db return db