Merge pull request 'feature-15' (#22) from feature-15 into main

Reviewed-on: #22
This commit is contained in:
nicolas.boufideline 2026-04-14 18:24:06 +00:00 committed by G3E Git Instance
commit a4a6be7feb
Signed by: G3E Git Instance
SSH key fingerprint: SHA256:7qPkHsv5cK9DqRLWKVhd6yvG6rpDxbWny9r8CMChJb0
16 changed files with 515 additions and 31 deletions

View file

@ -39,6 +39,7 @@ jobs:
- agent - agent
- vpc - vpc
- dhcp - dhcp
- subnet
uses: ./.forgejo/workflows/build.yml uses: ./.forgejo/workflows/build.yml
with: with:
tag: ${{ needs.set-release-target.outputs.release_cible }} tag: ${{ needs.set-release-target.outputs.release_cible }}

View file

@ -35,15 +35,19 @@ func main() {
defer db.Close() defer db.Close()
if *start { if *start {
metadata.StartMetadata(metadata.NoCloudConfig{ if err := metadata.StartMetadata(metadata.NoCloudConfig{
VpcName: *vpc, VpcName: *vpc,
Name: *vm_name, Name: *vm_name,
BindIP: *bind_ip, BindIP: *bind_ip,
BindPort: *bind_port, BindPort: *bind_port,
Password: *password, Password: *password,
SSHKEY: *ssh_key, SSHKEY: *ssh_key,
}, db, *dryrun) }, db, *dryrun); err != nil {
fmt.Println(err)
}
} else if *stop { } else if *stop {
metadata.StopMetadata(*vm_name, db, *dryrun) if err := metadata.StopMetadata(*vm_name, db, *dryrun); err != nil {
fmt.Println(err)
}
} }
} }

93
cmd/subnet/main.go Normal file
View file

@ -0,0 +1,93 @@
package main
import (
"flag"
"fmt"
"os"
configuration "git.g3e.fr/syonad/two/internal/config/agent"
"git.g3e.fr/syonad/two/internal/subnet"
"git.g3e.fr/syonad/two/pkg/db/kv"
"github.com/dgraph-io/badger/v4"
)
var (
name = flag.String("name", "", "Subnet name (ex: sn-00001)")
vpcName = flag.String("vpc", "", "VPC name")
vxlanID = flag.String("vxlan-id", "", "VXLAN ID")
localIP = flag.String("local-ip", "", "Local VTEP IP")
gatewayIP = flag.String("gateway-ip", "", "Gateway IP")
cidr = flag.String("cidr", "", "Subnet CIDR (ex: 10.10.10.0/24)")
action = flag.String("action", "", "Action à effectuer")
conf_file = flag.String("conf", "/etc/two/agent.yml", "Configuration file")
)
var DB *badger.DB
func main() {
flag.Parse()
conf, err := configuration.LoadConfig(*conf_file)
if err != nil {
fmt.Println(err)
os.Exit(1)
}
DB = kv.InitDB(kv.Config{
Path: conf.Database.Path,
}, false)
defer DB.Close()
switch *action {
case "create":
if *name == "" || *vpcName == "" || *vxlanID == "" || *localIP == "" || *gatewayIP == "" || *cidr == "" {
fmt.Println("create requires: -name -vpc -vxlan-id -local-ip -gateway-ip -cidr")
os.Exit(1)
}
kv.AddInDB(DB, "subnet/"+*name+"/state", "creating")
kv.AddInDB(DB, "subnet/"+*name+"/vpc", *vpcName)
kv.AddInDB(DB, "subnet/"+*name+"/vxlan_id", *vxlanID)
kv.AddInDB(DB, "subnet/"+*name+"/local_ip", *localIP)
kv.AddInDB(DB, "subnet/"+*name+"/gateway_ip", *gatewayIP)
kv.AddInDB(DB, "subnet/"+*name+"/cidr", *cidr)
if err := subnet.CreateSubnet(DB, *name); err != nil {
fmt.Println(err)
os.Exit(1)
}
case "delete":
if *name == "" {
fmt.Println("delete requires: -name")
os.Exit(1)
}
kv.AddInDB(DB, "subnet/"+*name+"/state", "deleting")
if err := subnet.DeleteSubnet(DB, *name); err != nil {
fmt.Println(err)
os.Exit(1)
}
if state, err := kv.GetFromDB(DB, "subnet/"+*name+"/state"); err != nil {
fmt.Println(err)
os.Exit(1)
} else if state == "deleted" {
kv.DeleteInDB(DB, "subnet/"+*name)
}
case "check":
if *name == "" {
fmt.Println("check requires: -name")
os.Exit(1)
}
if state, err := kv.GetFromDB(DB, "subnet/"+*name+"/state"); err != nil {
os.Exit(1)
} else if state != "created" {
os.Exit(1)
}
default:
fmt.Printf("Available commands:\n - create\n - delete\n - check\n")
os.Exit(1)
}
os.Exit(0)
}

View file

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

21
internal/netif/addr.go Normal file
View file

@ -0,0 +1,21 @@
package netif
import (
"net"
"github.com/vishvananda/netlink"
)
func AddrAdd(iface string, ip net.IP) error {
link, err := netlink.LinkByName(iface)
if err != nil {
return err
}
addr := &netlink.Addr{
IPNet: &net.IPNet{
IP: ip,
Mask: net.CIDRMask(32, 32),
},
}
return netlink.AddrAdd(link, addr)
}

View file

@ -0,0 +1,21 @@
//go:build linux
package netif
import (
"net"
"github.com/vishvananda/netlink"
)
func RouteAdd(iface string, subnet *net.IPNet) error {
link, err := netlink.LinkByName(iface)
if err != nil {
return err
}
return netlink.RouteAdd(&netlink.Route{
LinkIndex: link.Attrs().Index,
Dst: subnet,
Scope: netlink.SCOPE_LINK,
})
}

View file

@ -0,0 +1,9 @@
//go:build !linux
package netif
import "net"
func RouteAdd(_ string, _ *net.IPNet) error {
return nil
}

20
internal/netif/vxlan.go Normal file
View file

@ -0,0 +1,20 @@
package netif
import (
"net"
"github.com/vishvananda/netlink"
)
func CreateVxlan(name string, vxlanID int, localIP net.IP) error {
vxlan := &netlink.Vxlan{
LinkAttrs: netlink.LinkAttrs{
Name: name,
},
VxlanId: vxlanID,
Port: 4789,
SrcAddr: localIP,
Learning: false,
}
return netlink.LinkAdd(vxlan)
}

View file

@ -2,6 +2,6 @@
package netns package netns
func call(name string, fn func() error) error { func call(_ string, fn func() error) error {
return fn() return fn()
} }

View file

@ -3,12 +3,17 @@
package netns package netns
import ( import (
"fmt"
"os" "os"
"runtime"
"golang.org/x/sys/unix" "golang.org/x/sys/unix"
) )
func create(name string) error { func create(name string) error {
runtime.LockOSThread()
defer runtime.UnlockOSThread()
base := "/var/run/netns" base := "/var/run/netns"
path := base + "/" + name path := base + "/" + name
@ -16,6 +21,12 @@ func create(name string) error {
return err return err
} }
// si le fichier existe déjà, le démonter d'abord
if _, err := os.Stat(path); err == nil {
unix.Unmount(path, unix.MNT_DETACH)
os.Remove(path)
}
// fichier cible // fichier cible
f, err := os.Create(path) f, err := os.Create(path)
if err != nil { if err != nil {
@ -35,9 +46,12 @@ func create(name string) error {
return err return err
} }
// bind mount du netns courant vers /var/run/netns/<name> // bind mount du netns du thread courant vers /var/run/netns/<name>
// /proc/self/ns/net pointe vers le ns du processus (thread principal),
// pas du thread courant — il faut utiliser le tid explicitement
threadNsPath := fmt.Sprintf("/proc/self/task/%d/ns/net", unix.Gettid())
if err := unix.Mount( if err := unix.Mount(
"/proc/self/ns/net", threadNsPath,
path, path,
"", "",
unix.MS_BIND, unix.MS_BIND,

186
internal/subnet/create.go Normal file
View file

@ -0,0 +1,186 @@
package subnet
import (
"fmt"
"net"
"os/exec"
"strconv"
"strings"
"git.g3e.fr/syonad/two/internal/dhcp"
"git.g3e.fr/syonad/two/internal/netif"
"git.g3e.fr/syonad/two/internal/netns"
"git.g3e.fr/syonad/two/pkg/db/kv"
"git.g3e.fr/syonad/two/pkg/systemd"
"github.com/dgraph-io/badger/v4"
)
func CreateSubnet(db *badger.DB, subnetName string) error {
state, err := kv.GetFromDB(db, "subnet/"+subnetName+"/state")
if err != nil {
return err
}
if state != "creating" {
return nil
}
// lecture des paramètres depuis la DB
vpcName, err := kv.GetFromDB(db, "subnet/"+subnetName+"/vpc")
if err != nil {
return fmt.Errorf("get vpc: %w", err)
}
vxlanIDStr, err := kv.GetFromDB(db, "subnet/"+subnetName+"/vxlan_id")
if err != nil {
return fmt.Errorf("get vxlan_id: %w", err)
}
vxlanID, err := strconv.Atoi(vxlanIDStr)
if err != nil {
return fmt.Errorf("parse vxlan_id: %w", err)
}
localIPStr, err := kv.GetFromDB(db, "subnet/"+subnetName+"/local_ip")
if err != nil {
return fmt.Errorf("get local_ip: %w", err)
}
localIP := net.ParseIP(localIPStr)
if localIP == nil {
return fmt.Errorf("invalid local_ip: %s", localIPStr)
}
gatewayIPStr, err := kv.GetFromDB(db, "subnet/"+subnetName+"/gateway_ip")
if err != nil {
return fmt.Errorf("get gateway_ip: %w", err)
}
gatewayIP := net.ParseIP(gatewayIPStr)
if gatewayIP == nil {
return fmt.Errorf("invalid gateway_ip: %s", gatewayIPStr)
}
cidr, err := kv.GetFromDB(db, "subnet/"+subnetName+"/cidr")
if err != nil {
return fmt.Errorf("get cidr: %w", err)
}
_, subnet, err := net.ParseCIDR(cidr)
if err != nil {
return fmt.Errorf("parse cidr: %w", err)
}
// subnet_id = partie après le premier '-' (ex: "sn-00001" -> "00001")
subnetID := strings.SplitN(subnetName, "-", 2)[1]
bridge := "br-" + subnetID
vxlanIface := fmt.Sprintf("vxlan-%d", vxlanID)
// veth pair
if err := netif.CreateVethToNetns("v-"+subnetID+"-e", "v-"+subnetID+"-i", "/var/run/netns/"+vpcName, 1500); err != nil {
return fmt.Errorf("create veth: %w", err)
}
// bridge dans le root netns
if err := netif.CreateBridge(bridge, 1500); err != nil {
return fmt.Errorf("create bridge: %w", err)
}
// bridge dans le netns VPC
if err := netns.Call(vpcName, func() error {
return netif.CreateBridge(bridge, 1500)
}); err != nil {
return fmt.Errorf("create bridge in netns: %w", err)
}
// vxlan
if err := netif.CreateVxlan(vxlanIface, vxlanID, localIP); err != nil {
return fmt.Errorf("create vxlan: %w", err)
}
// ajout des interfaces dans les bridges
if err := netif.BridgeSetMaster("v-"+subnetID+"-e", bridge); err != nil {
return fmt.Errorf("add veth-e to bridge: %w", err)
}
if err := netns.Call(vpcName, func() error {
return netif.BridgeSetMaster("v-"+subnetID+"-i", bridge)
}); err != nil {
return fmt.Errorf("add veth-i to bridge in netns: %w", err)
}
if err := netif.BridgeSetMaster(vxlanIface, bridge); err != nil {
return fmt.Errorf("add vxlan to bridge: %w", err)
}
// montée des interfaces dans le root netns
for _, iface := range []string{"v-" + subnetID + "-e", vxlanIface, bridge} {
if err := netif.LinkSetUp(iface); err != nil {
return fmt.Errorf("set up %s: %w", iface, err)
}
}
// montée des interfaces dans le netns VPC
if err := netns.Call(vpcName, func() error {
for _, iface := range []string{"v-" + subnetID + "-i", bridge} {
if err := netif.LinkSetUp(iface); err != nil {
return fmt.Errorf("set up %s: %w", iface, err)
}
}
return nil
}); err != nil {
return fmt.Errorf("set up interfaces in netns: %w", err)
}
// IP gateway (/32) sur le bridge interne
if err := netns.Call(vpcName, func() error {
return netif.AddrAdd(bridge, gatewayIP)
}); err != nil {
return fmt.Errorf("add addr to bridge in netns: %w", err)
}
// route subnet (scope link) dans le netns VPC
if err := netns.Call(vpcName, func() error {
return netif.RouteAdd(bridge, subnet)
}); err != nil {
return fmt.Errorf("add route in netns: %w", err)
}
// ebtables : drop ARP Request vers la gateway sur ce bridge
if err := exec.Command("ebtables", "-A", "FORWARD",
"--out-interface", bridge,
"-p", "arp",
"--arp-op", "Request",
"--arp-ip-dst", gatewayIP.String(),
"-j", "DROP").Run(); err != nil {
return fmt.Errorf("ebtables arp rule: %w", err)
}
// ebtables : drop trafic DHCP sur ce bridge
if err := exec.Command("ebtables", "-A", "FORWARD",
"--out-interface", bridge,
"-p", "IPv4",
"--ip-protocol", "udp",
"--ip-source-port", "67:68",
"--ip-destination-port", "67:68",
"-j", "DROP").Run(); err != nil {
return fmt.Errorf("ebtables dhcp rule: %w", err)
}
// génération de la config dnsmasq et démarrage du service
conf := dhcp.Config{
Network: subnet,
Gateway: gatewayIP,
Name: vpcName + "_" + bridge,
ConfDir: "/etc/dnsmasq.d",
}
if _, err := dhcp.GenerateConfig(conf); err != nil {
return fmt.Errorf("generate dhcp config: %w", err)
}
svc, err := systemd.New()
if err != nil {
return fmt.Errorf("connect to systemd: %w", err)
}
defer svc.Close()
if err := svc.Start("dnsmasq@" + conf.Name + ".service"); err != nil {
return fmt.Errorf("start dnsmasq: %w", err)
}
return kv.AddInDB(db, "subnet/"+subnetName+"/state", "created")
}

95
internal/subnet/delete.go Normal file
View file

@ -0,0 +1,95 @@
package subnet
import (
"fmt"
"os"
"os/exec"
"strings"
"git.g3e.fr/syonad/two/internal/netif"
"git.g3e.fr/syonad/two/internal/netns"
"git.g3e.fr/syonad/two/pkg/db/kv"
"git.g3e.fr/syonad/two/pkg/systemd"
"github.com/dgraph-io/badger/v4"
)
func DeleteSubnet(db *badger.DB, subnetName string) error {
state, err := kv.GetFromDB(db, "subnet/"+subnetName+"/state")
if err != nil {
return err
}
if state != "deleting" {
return nil
}
vpcName, err := kv.GetFromDB(db, "subnet/"+subnetName+"/vpc")
if err != nil {
return fmt.Errorf("get vpc: %w", err)
}
vxlanIDStr, err := kv.GetFromDB(db, "subnet/"+subnetName+"/vxlan_id")
if err != nil {
return fmt.Errorf("get vxlan_id: %w", err)
}
subnetID := strings.SplitN(subnetName, "-", 2)[1]
bridge := "br-" + subnetID
vxlanIface := "vxlan-" + vxlanIDStr
// arrêt du service dnsmasq
svc, err := systemd.New()
if err != nil {
return fmt.Errorf("connect to systemd: %w", err)
}
defer svc.Close()
svcName := "dnsmasq@" + vpcName + "_" + bridge + ".service"
if err := svc.Stop(svcName); err != nil {
return fmt.Errorf("stop dnsmasq: %w", err)
}
// suppression de la config dnsmasq
if err := os.Remove("/etc/dnsmasq.d/" + vpcName + "_" + bridge + ".conf"); err != nil && !os.IsNotExist(err) {
return fmt.Errorf("remove dnsmasq config: %w", err)
}
// suppression des règles ebtables
exec.Command("ebtables", "-D", "FORWARD",
"--out-interface", bridge,
"-p", "arp",
"--arp-op", "Request",
"-j", "DROP").Run()
exec.Command("ebtables", "-D", "FORWARD",
"--out-interface", bridge,
"-p", "IPv4",
"--ip-protocol", "udp",
"--ip-source-port", "67:68",
"--ip-destination-port", "67:68",
"-j", "DROP").Run()
// suppression du bridge dans le netns VPC
if err := netns.Call(vpcName, func() error {
return netif.DeleteLink(bridge)
}); err != nil {
return fmt.Errorf("delete bridge in netns: %w", err)
}
// suppression du vxlan
if err := netif.DeleteLink(vxlanIface); err != nil {
return fmt.Errorf("delete vxlan: %w", err)
}
// suppression du veth pair (supprime les deux côtés)
if err := netif.DeleteLink("v-" + subnetID + "-e"); err != nil {
return fmt.Errorf("delete veth: %w", err)
}
// suppression du bridge dans le root netns
if err := netif.DeleteLink(bridge); err != nil {
return fmt.Errorf("delete bridge: %w", err)
}
return kv.AddInDB(db, "subnet/"+subnetName+"/state", "deleted")
}

View file

@ -22,7 +22,7 @@ func CreateVPC(db *badger.DB, name string) error {
} }
// create veth public for this netns // create veth public for this netns
if err := netif.CreateVethToNetns("veth"+name+"ext", "vethpublicint", "/var/run/netns/"+name, 9000); err != nil { if err := netif.CreateVethToNetns("vp-"+name+"-e", "vp-public-i", "/var/run/netns/"+name, 9000); err != nil {
return err return err
} }
@ -34,24 +34,24 @@ func CreateVPC(db *badger.DB, name string) error {
} }
// set veth to ext public bridge // set veth to ext public bridge
if err := netif.BridgeSetMaster("veth"+name+"ext", "br-public"); err != nil { if err := netif.BridgeSetMaster("vp-"+name+"-e", "br-public"); err != nil {
return err return err
} }
// set veth to int public bridge // set veth to int public bridge
if err := netns.Call(name, func() error { if err := netns.Call(name, func() error {
return netif.BridgeSetMaster("vethpublicint", "br-public") return netif.BridgeSetMaster("vp-public-i", "br-public")
}); err != nil { }); err != nil {
return err return err
} }
// set set ext veth up // set set ext veth up
if err := netif.LinkSetUp("veth" + name + "ext"); err != nil { if err := netif.LinkSetUp("vp-" + name + "-e"); err != nil {
return nil return err
} }
// set set int veth up // set set int veth up
if err := netns.Call(name, func() error { if err := netns.Call(name, func() error {
return netif.LinkSetUp("vethpublicint") return netif.LinkSetUp("vp-public-i")
}); err != nil { }); err != nil {
return err return err
} }

View file

@ -12,7 +12,7 @@ func DeleteVPC(db *badger.DB, name string) error {
if state, err := kv.GetFromDB(db, "vpc/"+name+"/state"); err != nil { if state, err := kv.GetFromDB(db, "vpc/"+name+"/state"); err != nil {
return err return err
} else if state == "deleting" { } else if state == "deleting" {
if err := netif.DeleteLink(name + "-ext"); err != nil { if err := netif.DeleteLink("vp-" + name + "-e"); err != nil {
return err return err
} }

View file

@ -6,7 +6,8 @@ import (
func InitDB(conf Config, readonly bool) *badger.DB { func InitDB(conf Config, readonly bool) *badger.DB {
opts := badger.DefaultOptions(conf.Path). opts := badger.DefaultOptions(conf.Path).
WithReadOnly(readonly) WithReadOnly(readonly).
WithBypassLockGuard(readonly)
opts.Logger = nil opts.Logger = nil
opts.ValueLogFileSize = 10 << 20 // 10 Mo par fichier vlog opts.ValueLogFileSize = 10 << 20 // 10 Mo par fichier vlog
opts.NumMemtables = 1 opts.NumMemtables = 1

View file

@ -11,6 +11,7 @@ import (
const ( const (
defaultTimeout = 5 * time.Second defaultTimeout = 5 * time.Second
jobTimeout = 30 * time.Second
jobMode = "replace" jobMode = "replace"
) )
@ -28,10 +29,7 @@ type ServiceStatus struct {
// New crée une connexion D-Bus systemd (scope système) // New crée une connexion D-Bus systemd (scope système)
func New() (*Manager, error) { func New() (*Manager, error) {
ctx, cancel := context.WithTimeout(context.Background(), defaultTimeout) conn, err := dbus.NewSystemConnectionContext(context.Background())
defer cancel()
conn, err := dbus.NewSystemConnectionContext(ctx)
if err != nil { if err != nil {
return nil, err return nil, err
} }
@ -57,17 +55,17 @@ func (m *Manager) Stop(service string) error {
} }
func (m *Manager) job(method, service string) error { func (m *Manager) job(method, service string) error {
ctx, cancel := context.WithTimeout(context.Background(), defaultTimeout) callCtx, callCancel := context.WithTimeout(context.Background(), defaultTimeout)
defer cancel() defer callCancel()
ch := make(chan string, 1) ch := make(chan string, 1)
var err error var err error
switch method { switch method {
case "StartUnit": case "StartUnit":
_, err = m.conn.StartUnitContext(ctx, service, jobMode, ch) _, err = m.conn.StartUnitContext(callCtx, service, jobMode, ch)
case "StopUnit": case "StopUnit":
_, err = m.conn.StopUnitContext(ctx, service, jobMode, ch) _, err = m.conn.StopUnitContext(callCtx, service, jobMode, ch)
default: default:
return errors.New("unsupported job method") return errors.New("unsupported job method")
} }
@ -76,9 +74,16 @@ func (m *Manager) job(method, service string) error {
return err return err
} }
result := <-ch waitCtx, waitCancel := context.WithTimeout(context.Background(), jobTimeout)
if result != "done" { defer waitCancel()
return fmt.Errorf("%s %s failed: %s", method, service, result)
select {
case result := <-ch:
if result != "done" {
return fmt.Errorf("%s %s failed: %s", method, service, result)
}
case <-waitCtx.Done():
return fmt.Errorf("%s %s timed out after %s", method, service, jobTimeout)
} }
return nil return nil