Compare commits

..

No commits in common. "25aebcc123c0718104c909078cd482c1113c42e1" and "374bb34695b67d12dcf8c5c4b7dd46b1a9d0a7dc" have entirely different histories.

13 changed files with 66 additions and 84 deletions

View file

@ -127,12 +127,6 @@ paths:
application/json: application/json:
schema: schema:
$ref: "#/components/schemas/Subnet" $ref: "#/components/schemas/Subnet"
"400":
description: Missing required field or unknown iface_type
content:
application/json:
schema:
$ref: "#/components/schemas/Error"
"409": "409":
description: Subnet already exists description: Subnet already exists
content: content:
@ -209,16 +203,15 @@ components:
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
pattern: '^vp-.+' example: vpc1
example: vp-00001
VPC: VPC:
type: object type: object
properties: properties:
name: name:
type: string type: string
example: vp-00001 example: vpc1
state: state:
type: string type: string
enum: [creating, created, deleting, deleted] enum: [creating, created, deleting, deleted]
@ -226,7 +219,7 @@ components:
SubnetCreateRequest: SubnetCreateRequest:
type: object type: object
required: [name, vpc, vxlan_id, gateway_ip, cidr] required: [name, vpc, vxlan_id, local_ip, gateway_ip, cidr]
properties: properties:
name: name:
type: string type: string
@ -240,10 +233,11 @@ components:
type: integer type: integer
description: VXLAN VNI identifier description: VXLAN VNI identifier
example: 100 example: 100
iface_type: local_ip:
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. format: ipv4
example: vms description: Local VTEP IP address
example: "10.0.0.5"
gateway_ip: gateway_ip:
type: string type: string
format: ipv4 format: ipv4
@ -270,10 +264,9 @@ components:
vxlan_id: vxlan_id:
type: integer type: integer
example: 100 example: 100
local_iface: local_ip:
type: string type: string
description: Resolved interface name example: "10.0.0.5"
example: br-000000
gateway_ip: gateway_ip:
type: string type: string
example: "10.10.10.1" example: "10.10.10.1"

View file

@ -36,7 +36,7 @@ func main() {
apiAddr := fmt.Sprintf("%s:%d", cfg.Api.Address, cfg.Api.Port) apiAddr := fmt.Sprintf("%s:%d", cfg.Api.Address, cfg.Api.Port)
promAddr := fmt.Sprintf("%s:%d", cfg.Prometheus.Address, cfg.Prometheus.Port) promAddr := fmt.Sprintf("%s:%d", cfg.Prometheus.Address, cfg.Prometheus.Port)
d := dispatcher.New(q, db, cfg) d := dispatcher.New(q, db)
go agentapi.New(d, db).Start(apiAddr) go agentapi.New(d, db).Start(apiAddr)
go promserver.Start(promAddr, registry) go promserver.Start(promAddr, registry)

View file

@ -1,9 +1,2 @@
database: database:
path: "/var/lib/two/data/" path: "/var/lib/two/data/"
default_interface: br-000000
interfaces:
vms: br-000000
internet: br-000000
admin: br-000000

View file

@ -13,7 +13,7 @@ type SubnetCreateRequest struct {
Name string `json:"name"` Name string `json:"name"`
VPC string `json:"vpc"` VPC string `json:"vpc"`
VxlanID int `json:"vxlan_id"` VxlanID int `json:"vxlan_id"`
IfaceType string `json:"iface_type"` LocalIP string `json:"local_ip"`
GatewayIP string `json:"gateway_ip"` GatewayIP string `json:"gateway_ip"`
CIDR string `json:"cidr"` CIDR string `json:"cidr"`
} }
@ -23,7 +23,7 @@ type Subnet struct {
State string `json:"state"` State string `json:"state"`
VPC string `json:"vpc"` VPC string `json:"vpc"`
VxlanID int `json:"vxlan_id"` VxlanID int `json:"vxlan_id"`
LocalIface string `json:"local_iface"` LocalIP string `json:"local_ip"`
GatewayIP string `json:"gateway_ip"` GatewayIP string `json:"gateway_ip"`
CIDR string `json:"cidr"` CIDR string `json:"cidr"`
} }

View file

@ -31,16 +31,16 @@ 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.IfaceType == "" || req.GatewayIP == "" || req.CIDR == "" { if req.Name == "" || req.VPC == "" || req.LocalIP == "" || req.GatewayIP == "" || req.CIDR == "" {
w.WriteHeader(http.StatusBadRequest) w.WriteHeader(http.StatusBadRequest)
json.NewEncoder(w).Encode(ErrorResponse{Error: "name, vpc, iface_type, gateway_ip and cidr are required"}) json.NewEncoder(w).Encode(ErrorResponse{Error: "name, vpc, local_ip, gateway_ip and cidr are required"})
return return
} }
s.dispatcher.Dispatch(dispatcher.CreateSubnetCommand{ s.dispatcher.Dispatch(dispatcher.CreateSubnetCommand{
Name: req.Name, Name: req.Name,
VPC: req.VPC, VPC: req.VPC,
VxlanID: req.VxlanID, VxlanID: req.VxlanID,
IfaceType: req.IfaceType, LocalIP: req.LocalIP,
GatewayIP: req.GatewayIP, GatewayIP: req.GatewayIP,
CIDR: req.CIDR, CIDR: req.CIDR,
}) })
@ -50,7 +50,9 @@ func (s *Server) postSubnet(w http.ResponseWriter, r *http.Request) {
State: "creating", State: "creating",
VPC: req.VPC, VPC: req.VPC,
VxlanID: req.VxlanID, VxlanID: req.VxlanID,
LocalIP: req.LocalIP,
GatewayIP: req.GatewayIP, GatewayIP: req.GatewayIP,
CIDR: req.CIDR, CIDR: req.CIDR,
}) })
} }

View file

@ -20,8 +20,6 @@ type Config struct {
Count int `mapstructure:"count"` Count int `mapstructure:"count"`
BufferSize int `mapstructure:"buffer_size"` BufferSize int `mapstructure:"buffer_size"`
} `mapstructure:"worker"` } `mapstructure:"worker"`
DefaultInterface string `mapstructure:"default_interface"`
Interfaces map[string]string `mapstructure:"interfaces"`
} }
func LoadConfig(path string) (*Config, error) { func LoadConfig(path string) (*Config, error) {
@ -36,7 +34,6 @@ func LoadConfig(path string) (*Config, error) {
v.SetDefault("prometheus.port", 9090) v.SetDefault("prometheus.port", 9090)
v.SetDefault("worker.count", 4) v.SetDefault("worker.count", 4)
v.SetDefault("worker.buffer_size", 100) v.SetDefault("worker.buffer_size", 100)
v.SetDefault("default_interface", "br-000000")
v.ReadInConfig() v.ReadInConfig()

View file

@ -3,28 +3,26 @@ package dispatcher
import ( import (
"log" "log"
configuration "git.g3e.fr/syonad/two/internal/config/agent"
"git.g3e.fr/syonad/two/pkg/worker" "git.g3e.fr/syonad/two/pkg/worker"
"github.com/dgraph-io/badger/v4" "github.com/dgraph-io/badger/v4"
) )
type Command interface { type Command interface {
Execute(db *badger.DB, cfg *configuration.Config) error Execute(db *badger.DB) error
} }
type Dispatcher struct { type Dispatcher struct {
queue *worker.Queue queue *worker.Queue
db *badger.DB db *badger.DB
cfg *configuration.Config
} }
func New(queue *worker.Queue, db *badger.DB, cfg *configuration.Config) *Dispatcher { func New(queue *worker.Queue, db *badger.DB) *Dispatcher {
return &Dispatcher{queue: queue, db: db, cfg: cfg} return &Dispatcher{queue: queue, db: db}
} }
func (d *Dispatcher) Dispatch(cmd Command) { func (d *Dispatcher) Dispatch(cmd Command) {
d.queue.Submit(func() { d.queue.Submit(func() {
if err := cmd.Execute(d.db, d.cfg); err != nil { if err := cmd.Execute(d.db); err != nil {
log.Printf("command error (%T): %v", cmd, err) log.Printf("command error (%T): %v", cmd, err)
} }
}) })

View file

@ -5,7 +5,6 @@ import (
"os" "os"
"strconv" "strconv"
configuration "git.g3e.fr/syonad/two/internal/config/agent"
"git.g3e.fr/syonad/two/internal/subnet" "git.g3e.fr/syonad/two/internal/subnet"
"git.g3e.fr/syonad/two/pkg/db/kv" "git.g3e.fr/syonad/two/pkg/db/kv"
"github.com/dgraph-io/badger/v4" "github.com/dgraph-io/badger/v4"
@ -15,20 +14,16 @@ type CreateSubnetCommand struct {
Name string Name string
VPC string VPC string
VxlanID int VxlanID int
IfaceType string LocalIP string
GatewayIP string GatewayIP string
CIDR string CIDR string
} }
func (c CreateSubnetCommand) Execute(db *badger.DB, cfg *configuration.Config) error { func (c CreateSubnetCommand) Execute(db *badger.DB) error {
localIface, ok := cfg.Interfaces[c.IfaceType]
if !ok {
localIface = cfg.DefaultInterface
}
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+"/vxlan_id", strconv.Itoa(c.VxlanID))
kv.AddInDB(db, "subnet/"+c.Name+"/local_iface", localIface) kv.AddInDB(db, "subnet/"+c.Name+"/local_ip", c.LocalIP)
kv.AddInDB(db, "subnet/"+c.Name+"/gateway_ip", c.GatewayIP) kv.AddInDB(db, "subnet/"+c.Name+"/gateway_ip", c.GatewayIP)
kv.AddInDB(db, "subnet/"+c.Name+"/cidr", c.CIDR) kv.AddInDB(db, "subnet/"+c.Name+"/cidr", c.CIDR)
return subnet.CreateSubnet(db, c.Name) return subnet.CreateSubnet(db, c.Name)
@ -38,7 +33,7 @@ type DeleteSubnetCommand struct {
Name string Name string
} }
func (c DeleteSubnetCommand) Execute(db *badger.DB, _ *configuration.Config) error { func (c DeleteSubnetCommand) Execute(db *badger.DB) error {
kv.AddInDB(db, "subnet/"+c.Name+"/state", "deleting") kv.AddInDB(db, "subnet/"+c.Name+"/state", "deleting")
if err := subnet.DeleteSubnet(db, c.Name); err != nil { if err := subnet.DeleteSubnet(db, c.Name); err != nil {
fmt.Println(err) fmt.Println(err)

View file

@ -1,7 +1,6 @@
package dispatcher package dispatcher
import ( import (
configuration "git.g3e.fr/syonad/two/internal/config/agent"
"git.g3e.fr/syonad/two/internal/vpc" "git.g3e.fr/syonad/two/internal/vpc"
"git.g3e.fr/syonad/two/pkg/db/kv" "git.g3e.fr/syonad/two/pkg/db/kv"
"github.com/dgraph-io/badger/v4" "github.com/dgraph-io/badger/v4"
@ -11,7 +10,7 @@ type CreateVPCCommand struct {
Name string Name string
} }
func (c CreateVPCCommand) Execute(db *badger.DB, _ *configuration.Config) error { func (c CreateVPCCommand) Execute(db *badger.DB) error {
kv.AddInDB(db, "vpc/"+c.Name+"/state", "creating") kv.AddInDB(db, "vpc/"+c.Name+"/state", "creating")
return vpc.CreateVPC(db, c.Name) return vpc.CreateVPC(db, c.Name)
} }
@ -20,7 +19,7 @@ type DeleteVPCCommand struct {
Name string Name string
} }
func (c DeleteVPCCommand) Execute(db *badger.DB, _ *configuration.Config) error { func (c DeleteVPCCommand) Execute(db *badger.DB) error {
kv.AddInDB(db, "vpc/"+c.Name+"/state", "deleting") kv.AddInDB(db, "vpc/"+c.Name+"/state", "deleting")
if err := vpc.DeleteVPC(db, c.Name); err != nil { if err := vpc.DeleteVPC(db, c.Name); err != nil {
return err return err

View file

@ -1,21 +1,19 @@
package netif package netif
import ( import (
"net"
"github.com/vishvananda/netlink" "github.com/vishvananda/netlink"
) )
func CreateVxlan(name string, vxlanID int, localIface string) error { func CreateVxlan(name string, vxlanID int, localIP net.IP) error {
link, err := netlink.LinkByName(localIface)
if err != nil {
return err
}
vxlan := &netlink.Vxlan{ vxlan := &netlink.Vxlan{
LinkAttrs: netlink.LinkAttrs{ LinkAttrs: netlink.LinkAttrs{
Name: name, Name: name,
}, },
VxlanId: vxlanID, VxlanId: vxlanID,
Port: 4789, Port: 4789,
VtepDevIndex: link.Attrs().Index, SrcAddr: localIP,
Learning: false, Learning: false,
} }
return netlink.LinkAdd(vxlan) return netlink.LinkAdd(vxlan)

View file

@ -40,9 +40,13 @@ func CreateSubnet(db *badger.DB, subnetName string) error {
return fmt.Errorf("parse vxlan_id: %w", err) return fmt.Errorf("parse vxlan_id: %w", err)
} }
localIface, err := kv.GetFromDB(db, "subnet/"+subnetName+"/local_iface") localIPStr, err := kv.GetFromDB(db, "subnet/"+subnetName+"/local_ip")
if err != nil { if err != nil {
return fmt.Errorf("get local_iface: %w", err) 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") gatewayIPStr, err := kv.GetFromDB(db, "subnet/"+subnetName+"/gateway_ip")
@ -86,7 +90,7 @@ func CreateSubnet(db *badger.DB, subnetName string) error {
} }
// vxlan // vxlan
if err := netif.CreateVxlan(vxlanIface, vxlanID, localIface); err != nil { if err := netif.CreateVxlan(vxlanIface, vxlanID, localIP); err != nil {
return fmt.Errorf("create vxlan: %w", err) return fmt.Errorf("create vxlan: %w", err)
} }

View file

@ -1,8 +1,6 @@
package vpc package vpc
import ( import (
"strings"
"git.g3e.fr/syonad/two/internal/netif" "git.g3e.fr/syonad/two/internal/netif"
"git.g3e.fr/syonad/two/internal/netns" "git.g3e.fr/syonad/two/internal/netns"
"git.g3e.fr/syonad/two/pkg/db/kv" "git.g3e.fr/syonad/two/pkg/db/kv"
@ -11,40 +9,49 @@ import (
) )
func CreateVPC(db *badger.DB, name string) error { func CreateVPC(db *badger.DB, name string) error {
// missing
// search data in db
// change state in db
// create netns
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 == "creating" { } else if state == "creating" {
vpcID := strings.SplitN(name, "-", 2)[1]
if err := netns.Create(name); err != nil { if err := netns.Create(name); err != nil {
return err return err
} }
if err := netif.CreateVethToNetns("vp-"+vpcID+"-e", "vp-"+vpcID+"-i", "/var/run/netns/"+name, 9000); err != nil { // create veth public for this netns
if err := netif.CreateVethToNetns("vp-"+name+"-e", "vp-public-i", "/var/run/netns/"+name, 9000); err != nil {
return err return err
} }
// create public bridge in netns
if err := netns.Call(name, func() error { if err := netns.Call(name, func() error {
return netif.CreateBridge("br-public", 1500) return netif.CreateBridge("br-public", 1500)
}); err != nil { }); err != nil {
return err return err
} }
if err := netif.BridgeSetMaster("vp-"+vpcID+"-e", "br-public"); err != nil { // set veth to ext public bridge
if err := netif.BridgeSetMaster("vp-"+name+"-e", "br-public"); err != nil {
return err return err
} }
// set veth to int public bridge
if err := netns.Call(name, func() error { if err := netns.Call(name, func() error {
return netif.BridgeSetMaster("vp-"+vpcID+"-i", "br-public") return netif.BridgeSetMaster("vp-public-i", "br-public")
}); err != nil { }); err != nil {
return err return err
} }
if err := netif.LinkSetUp("vp-" + vpcID + "-e"); err != nil { // set set ext veth up
if err := netif.LinkSetUp("vp-" + name + "-e"); err != nil {
return err return err
} }
// set set int veth up
if err := netns.Call(name, func() error { if err := netns.Call(name, func() error {
return netif.LinkSetUp("vp-" + vpcID + "-i") return netif.LinkSetUp("vp-public-i")
}); err != nil { }); err != nil {
return err return err
} }

View file

@ -1,8 +1,6 @@
package vpc package vpc
import ( import (
"strings"
"git.g3e.fr/syonad/two/internal/netif" "git.g3e.fr/syonad/two/internal/netif"
"git.g3e.fr/syonad/two/internal/netns" "git.g3e.fr/syonad/two/internal/netns"
"git.g3e.fr/syonad/two/pkg/db/kv" "git.g3e.fr/syonad/two/pkg/db/kv"
@ -14,9 +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" {
vpcID := strings.SplitN(name, "-", 2)[1] if err := netif.DeleteLink("vp-" + name + "-e"); err != nil {
if err := netif.DeleteLink("vp-" + vpcID + "-e"); err != nil {
return err return err
} }