Compare commits
3 commits
19434b1848
...
648a64782c
| Author | SHA1 | Date | |
|---|---|---|---|
|
648a64782c |
|||
|
9ee361ae28 |
|||
|
1d86c45ba4 |
9 changed files with 194 additions and 22 deletions
|
|
@ -1,8 +1,35 @@
|
||||||
|
# Path to the Badger key-value database directory
|
||||||
database:
|
database:
|
||||||
path: "/var/lib/two/data/"
|
path: "/var/lib/two/data/"
|
||||||
|
|
||||||
|
# REST API server
|
||||||
|
api:
|
||||||
|
address: "0.0.0.0"
|
||||||
|
port: 8080
|
||||||
|
|
||||||
|
# Prometheus metrics server
|
||||||
|
prometheus:
|
||||||
|
address: "0.0.0.0"
|
||||||
|
port: 9090
|
||||||
|
|
||||||
|
# Worker pool that executes dispatched commands
|
||||||
|
worker:
|
||||||
|
# Number of concurrent worker goroutines
|
||||||
|
count: 4
|
||||||
|
# Maximum number of commands queued before Dispatch blocks
|
||||||
|
buffer_size: 100
|
||||||
|
|
||||||
|
# Timing for commands that wait on resource state transitions
|
||||||
|
dispatcher:
|
||||||
|
# How long (in seconds) to wait before giving up
|
||||||
|
timeout_seconds: 300
|
||||||
|
# Interval (in seconds) between each state check
|
||||||
|
poll_seconds: 2
|
||||||
|
|
||||||
|
# Bridge interface used when the requested iface_type is not found in interfaces
|
||||||
default_interface: br-000000
|
default_interface: br-000000
|
||||||
|
|
||||||
|
# Map of logical interface types to physical bridge names on this host
|
||||||
interfaces:
|
interfaces:
|
||||||
vms: br-000000
|
vms: br-000000
|
||||||
internet: br-000000
|
internet: br-000000
|
||||||
|
|
|
||||||
|
|
@ -61,8 +61,15 @@ func (s *Server) getSubnet(w http.ResponseWriter, _ *http.Request, name string)
|
||||||
json.NewEncoder(w).Encode(sub)
|
json.NewEncoder(w).Encode(sub)
|
||||||
}
|
}
|
||||||
|
|
||||||
func (s *Server) deleteSubnet(w http.ResponseWriter, r *http.Request, name string) {
|
func (s *Server) deleteSubnet(w http.ResponseWriter, _ *http.Request, name string) {
|
||||||
s.dispatcher.Dispatch(dispatcher.DeleteSubnetCommand{Name: name})
|
cmd := dispatcher.DeleteSubnetCommand{Name: name}
|
||||||
w.WriteHeader(http.StatusAccepted)
|
if err := s.dispatcher.Prepare(cmd); err != nil {
|
||||||
json.NewEncoder(w).Encode(Subnet{Name: name, State: "deleting"})
|
w.WriteHeader(http.StatusNotFound)
|
||||||
|
json.NewEncoder(w).Encode(ErrorResponse{Error: err.Error()})
|
||||||
|
return
|
||||||
|
}
|
||||||
|
s.dispatcher.Dispatch(cmd)
|
||||||
|
state, _ := kv.GetFromDB(s.db, "subnet/"+name+"/state")
|
||||||
|
w.WriteHeader(http.StatusAccepted)
|
||||||
|
json.NewEncoder(w).Encode(Subnet{Name: name, State: state})
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -74,21 +74,42 @@ func (s *Server) postSubnet(w http.ResponseWriter, r *http.Request) {
|
||||||
json.NewEncoder(w).Encode(ErrorResponse{Error: "name, vpc, iface_type, gateway_ip and cidr are required"})
|
json.NewEncoder(w).Encode(ErrorResponse{Error: "name, vpc, iface_type, gateway_ip and cidr are required"})
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
s.dispatcher.Dispatch(dispatcher.CreateSubnetCommand{
|
cmd := dispatcher.CreateSubnetCommand{
|
||||||
Name: req.Name,
|
Name: req.Name,
|
||||||
VPC: req.VPC,
|
VPC: req.VPC,
|
||||||
VxlanID: req.VxlanID,
|
VxlanID: req.VxlanID,
|
||||||
IfaceType: req.IfaceType,
|
IfaceType: req.IfaceType,
|
||||||
GatewayIP: req.GatewayIP,
|
GatewayIP: req.GatewayIP,
|
||||||
CIDR: req.CIDR,
|
CIDR: req.CIDR,
|
||||||
})
|
}
|
||||||
w.WriteHeader(http.StatusAccepted)
|
if err := s.dispatcher.Prepare(cmd); err != nil {
|
||||||
json.NewEncoder(w).Encode(Subnet{
|
w.WriteHeader(http.StatusConflict)
|
||||||
Name: req.Name,
|
json.NewEncoder(w).Encode(ErrorResponse{Error: err.Error()})
|
||||||
State: "creating",
|
return
|
||||||
VPC: req.VPC,
|
}
|
||||||
VxlanID: req.VxlanID,
|
s.dispatcher.Dispatch(cmd)
|
||||||
GatewayIP: req.GatewayIP,
|
entries, _ := kv.ListByPrefix(s.db, "subnet/"+req.Name+"/")
|
||||||
CIDR: req.CIDR,
|
sub := Subnet{Name: req.Name}
|
||||||
})
|
for key, value := range entries {
|
||||||
|
parts := strings.Split(key, "/")
|
||||||
|
if len(parts) != 3 {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
switch parts[2] {
|
||||||
|
case "state":
|
||||||
|
sub.State = value
|
||||||
|
case "vpc":
|
||||||
|
sub.VPC = value
|
||||||
|
case "vxlan_id":
|
||||||
|
sub.VxlanID, _ = strconv.Atoi(value)
|
||||||
|
case "local_iface":
|
||||||
|
sub.LocalIface = value
|
||||||
|
case "gateway_ip":
|
||||||
|
sub.GatewayIP = value
|
||||||
|
case "cidr":
|
||||||
|
sub.CIDR = value
|
||||||
|
}
|
||||||
|
}
|
||||||
|
w.WriteHeader(http.StatusAccepted)
|
||||||
|
json.NewEncoder(w).Encode(sub)
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -40,7 +40,14 @@ func (s *Server) getVpc(w http.ResponseWriter, _ *http.Request, name string) {
|
||||||
}
|
}
|
||||||
|
|
||||||
func (s *Server) deleteVpc(w http.ResponseWriter, _ *http.Request, name string) {
|
func (s *Server) deleteVpc(w http.ResponseWriter, _ *http.Request, name string) {
|
||||||
s.dispatcher.Dispatch(dispatcher.DeleteVPCCommand{Name: name})
|
cmd := dispatcher.DeleteVPCCommand{Name: name}
|
||||||
w.WriteHeader(http.StatusAccepted)
|
if err := s.dispatcher.Prepare(cmd); err != nil {
|
||||||
json.NewEncoder(w).Encode(VPC{Name: name, State: "deleting"})
|
w.WriteHeader(http.StatusNotFound)
|
||||||
|
json.NewEncoder(w).Encode(ErrorResponse{Error: err.Error()})
|
||||||
|
return
|
||||||
|
}
|
||||||
|
s.dispatcher.Dispatch(cmd)
|
||||||
|
state, _ := kv.GetFromDB(s.db, "vpc/"+name+"/state")
|
||||||
|
w.WriteHeader(http.StatusAccepted)
|
||||||
|
json.NewEncoder(w).Encode(VPC{Name: name, State: state})
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -62,7 +62,14 @@ 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
|
||||||
}
|
}
|
||||||
s.dispatcher.Dispatch(dispatcher.CreateVPCCommand{Name: req.Name})
|
cmd := dispatcher.CreateVPCCommand{Name: req.Name}
|
||||||
w.WriteHeader(http.StatusAccepted)
|
if err := s.dispatcher.Prepare(cmd); err != nil {
|
||||||
json.NewEncoder(w).Encode(VPC{Name: req.Name, State: "creating"})
|
w.WriteHeader(http.StatusConflict)
|
||||||
|
json.NewEncoder(w).Encode(ErrorResponse{Error: err.Error()})
|
||||||
|
return
|
||||||
|
}
|
||||||
|
s.dispatcher.Dispatch(cmd)
|
||||||
|
state, _ := kv.GetFromDB(s.db, "vpc/"+req.Name+"/state")
|
||||||
|
w.WriteHeader(http.StatusAccepted)
|
||||||
|
json.NewEncoder(w).Encode(VPC{Name: req.Name, State: state})
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -20,6 +20,10 @@ 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"`
|
||||||
|
Dispatcher struct {
|
||||||
|
TimeoutSeconds int `mapstructure:"timeout_seconds"`
|
||||||
|
PollSeconds int `mapstructure:"poll_seconds"`
|
||||||
|
} `mapstructure:"dispatcher"`
|
||||||
DefaultInterface string `mapstructure:"default_interface"`
|
DefaultInterface string `mapstructure:"default_interface"`
|
||||||
Interfaces map[string]string `mapstructure:"interfaces"`
|
Interfaces map[string]string `mapstructure:"interfaces"`
|
||||||
}
|
}
|
||||||
|
|
@ -36,6 +40,8 @@ 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("dispatcher.timeout_seconds", 300)
|
||||||
|
v.SetDefault("dispatcher.poll_seconds", 2)
|
||||||
v.SetDefault("default_interface", "br-000000")
|
v.SetDefault("default_interface", "br-000000")
|
||||||
|
|
||||||
v.ReadInConfig()
|
v.ReadInConfig()
|
||||||
|
|
|
||||||
|
|
@ -9,6 +9,7 @@ import (
|
||||||
)
|
)
|
||||||
|
|
||||||
type Command interface {
|
type Command interface {
|
||||||
|
Prepare(db *badger.DB, cfg *configuration.Config) error
|
||||||
Execute(db *badger.DB, cfg *configuration.Config) error
|
Execute(db *badger.DB, cfg *configuration.Config) error
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -22,6 +23,10 @@ func New(queue *worker.Queue, db *badger.DB, cfg *configuration.Config) *Dispatc
|
||||||
return &Dispatcher{queue: queue, db: db, cfg: cfg}
|
return &Dispatcher{queue: queue, db: db, cfg: cfg}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (d *Dispatcher) Prepare(cmd Command) error {
|
||||||
|
return cmd.Prepare(d.db, d.cfg)
|
||||||
|
}
|
||||||
|
|
||||||
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, d.cfg); err != nil {
|
||||||
|
|
|
||||||
|
|
@ -4,6 +4,7 @@ import (
|
||||||
"fmt"
|
"fmt"
|
||||||
"os"
|
"os"
|
||||||
"strconv"
|
"strconv"
|
||||||
|
"time"
|
||||||
|
|
||||||
configuration "git.g3e.fr/syonad/two/internal/config/agent"
|
configuration "git.g3e.fr/syonad/two/internal/config/agent"
|
||||||
"git.g3e.fr/syonad/two/internal/subnet"
|
"git.g3e.fr/syonad/two/internal/subnet"
|
||||||
|
|
@ -20,7 +21,17 @@ type CreateSubnetCommand struct {
|
||||||
CIDR string
|
CIDR string
|
||||||
}
|
}
|
||||||
|
|
||||||
func (c CreateSubnetCommand) Execute(db *badger.DB, cfg *configuration.Config) error {
|
func (c CreateSubnetCommand) Prepare(db *badger.DB, cfg *configuration.Config) error {
|
||||||
|
if _, err := kv.GetFromDB(db, "subnet/"+c.Name+"/state"); err == nil {
|
||||||
|
return fmt.Errorf("subnet %q already exists", c.Name)
|
||||||
|
}
|
||||||
|
vpcState, err := kv.GetFromDB(db, "vpc/"+c.VPC+"/state")
|
||||||
|
if err != nil {
|
||||||
|
return fmt.Errorf("vpc %q not found", c.VPC)
|
||||||
|
}
|
||||||
|
if vpcState == "deleting" || vpcState == "deleted" {
|
||||||
|
return fmt.Errorf("vpc %q is %s", c.VPC, vpcState)
|
||||||
|
}
|
||||||
localIface, ok := cfg.Interfaces[c.IfaceType]
|
localIface, ok := cfg.Interfaces[c.IfaceType]
|
||||||
if !ok {
|
if !ok {
|
||||||
localIface = cfg.DefaultInterface
|
localIface = cfg.DefaultInterface
|
||||||
|
|
@ -31,6 +42,25 @@ func (c CreateSubnetCommand) Execute(db *badger.DB, cfg *configuration.Config) e
|
||||||
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+"/gateway_ip", c.GatewayIP)
|
||||||
kv.AddInDB(db, "subnet/"+c.Name+"/cidr", c.CIDR)
|
kv.AddInDB(db, "subnet/"+c.Name+"/cidr", c.CIDR)
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func (c CreateSubnetCommand) Execute(db *badger.DB, cfg *configuration.Config) error {
|
||||||
|
timeout := time.After(time.Duration(cfg.Dispatcher.TimeoutSeconds) * time.Second)
|
||||||
|
for {
|
||||||
|
state, err := kv.GetFromDB(db, "vpc/"+c.VPC+"/state")
|
||||||
|
if err != nil {
|
||||||
|
return fmt.Errorf("vpc %q not found while waiting", c.VPC)
|
||||||
|
}
|
||||||
|
if state == "created" {
|
||||||
|
break
|
||||||
|
}
|
||||||
|
select {
|
||||||
|
case <-timeout:
|
||||||
|
return fmt.Errorf("timed out waiting for vpc %q to be created", c.VPC)
|
||||||
|
case <-time.After(time.Duration(cfg.Dispatcher.PollSeconds) * time.Second):
|
||||||
|
}
|
||||||
|
}
|
||||||
return subnet.CreateSubnet(db, c.Name)
|
return subnet.CreateSubnet(db, c.Name)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -38,8 +68,14 @@ type DeleteSubnetCommand struct {
|
||||||
Name string
|
Name string
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (c DeleteSubnetCommand) Prepare(db *badger.DB, _ *configuration.Config) error {
|
||||||
|
if _, err := kv.GetFromDB(db, "subnet/"+c.Name+"/state"); err != nil {
|
||||||
|
return fmt.Errorf("subnet %q not found", c.Name)
|
||||||
|
}
|
||||||
|
return kv.AddInDB(db, "subnet/"+c.Name+"/state", "deleting")
|
||||||
|
}
|
||||||
|
|
||||||
func (c DeleteSubnetCommand) Execute(db *badger.DB, _ *configuration.Config) error {
|
func (c DeleteSubnetCommand) Execute(db *badger.DB, _ *configuration.Config) error {
|
||||||
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)
|
||||||
os.Exit(1)
|
os.Exit(1)
|
||||||
|
|
|
||||||
|
|
@ -1,6 +1,10 @@
|
||||||
package dispatcher
|
package dispatcher
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"fmt"
|
||||||
|
"strings"
|
||||||
|
"time"
|
||||||
|
|
||||||
configuration "git.g3e.fr/syonad/two/internal/config/agent"
|
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"
|
||||||
|
|
@ -11,8 +15,14 @@ type CreateVPCCommand struct {
|
||||||
Name string
|
Name string
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (c CreateVPCCommand) Prepare(db *badger.DB, _ *configuration.Config) error {
|
||||||
|
if _, err := kv.GetFromDB(db, "vpc/"+c.Name+"/state"); err == nil {
|
||||||
|
return fmt.Errorf("vpc %q already exists", c.Name)
|
||||||
|
}
|
||||||
|
return kv.AddInDB(db, "vpc/"+c.Name+"/state", "creating")
|
||||||
|
}
|
||||||
|
|
||||||
func (c CreateVPCCommand) Execute(db *badger.DB, _ *configuration.Config) error {
|
func (c CreateVPCCommand) Execute(db *badger.DB, _ *configuration.Config) error {
|
||||||
kv.AddInDB(db, "vpc/"+c.Name+"/state", "creating")
|
|
||||||
return vpc.CreateVPC(db, c.Name)
|
return vpc.CreateVPC(db, c.Name)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -20,8 +30,54 @@ type DeleteVPCCommand struct {
|
||||||
Name string
|
Name string
|
||||||
}
|
}
|
||||||
|
|
||||||
func (c DeleteVPCCommand) Execute(db *badger.DB, _ *configuration.Config) error {
|
func (c DeleteVPCCommand) Prepare(db *badger.DB, _ *configuration.Config) error {
|
||||||
kv.AddInDB(db, "vpc/"+c.Name+"/state", "deleting")
|
if _, err := kv.GetFromDB(db, "vpc/"+c.Name+"/state"); err != nil {
|
||||||
|
return fmt.Errorf("vpc %q not found", c.Name)
|
||||||
|
}
|
||||||
|
entries, err := kv.ListByPrefix(db, "subnet/")
|
||||||
|
if err != nil {
|
||||||
|
return fmt.Errorf("failed to list subnets: %w", err)
|
||||||
|
}
|
||||||
|
for key, value := range entries {
|
||||||
|
if !strings.HasSuffix(key, "/vpc") || value != c.Name {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
subnetName := strings.Split(key, "/")[1]
|
||||||
|
state, err := kv.GetFromDB(db, "subnet/"+subnetName+"/state")
|
||||||
|
if err != nil || (state != "deleting" && state != "deleted") {
|
||||||
|
return fmt.Errorf("subnet %q must be deleted before deleting vpc %q", subnetName, c.Name)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return kv.AddInDB(db, "vpc/"+c.Name+"/state", "deleting")
|
||||||
|
}
|
||||||
|
|
||||||
|
func (c DeleteVPCCommand) Execute(db *badger.DB, cfg *configuration.Config) error {
|
||||||
|
timeout := time.After(time.Duration(cfg.Dispatcher.TimeoutSeconds) * time.Second)
|
||||||
|
for {
|
||||||
|
entries, err := kv.ListByPrefix(db, "subnet/")
|
||||||
|
if err != nil {
|
||||||
|
return fmt.Errorf("failed to list subnets: %w", err)
|
||||||
|
}
|
||||||
|
pending := false
|
||||||
|
for key, value := range entries {
|
||||||
|
if strings.HasSuffix(key, "/vpc") && value == c.Name {
|
||||||
|
subnetName := strings.Split(key, "/")[1]
|
||||||
|
state, _ := kv.GetFromDB(db, "subnet/"+subnetName+"/state")
|
||||||
|
if state == "deleting" {
|
||||||
|
pending = true
|
||||||
|
break
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if !pending {
|
||||||
|
break
|
||||||
|
}
|
||||||
|
select {
|
||||||
|
case <-timeout:
|
||||||
|
return fmt.Errorf("timed out waiting for subnets of vpc %q to be deleted", c.Name)
|
||||||
|
case <-time.After(time.Duration(cfg.Dispatcher.PollSeconds) * time.Second):
|
||||||
|
}
|
||||||
|
}
|
||||||
if err := vpc.DeleteVPC(db, c.Name); err != nil {
|
if err := vpc.DeleteVPC(db, c.Name); err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue