diff --git a/conf/agent/config.exemple.yml b/conf/agent/config.exemple.yml index 0a966dd..bc1d437 100644 --- a/conf/agent/config.exemple.yml +++ b/conf/agent/config.exemple.yml @@ -1,35 +1,8 @@ -# Path to the Badger key-value database directory database: 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 -# Map of logical interface types to physical bridge names on this host interfaces: vms: br-000000 internet: br-000000 diff --git a/internal/api/agent/subnet.go b/internal/api/agent/subnet.go index 556b7c3..61cbd6f 100644 --- a/internal/api/agent/subnet.go +++ b/internal/api/agent/subnet.go @@ -61,15 +61,8 @@ func (s *Server) getSubnet(w http.ResponseWriter, _ *http.Request, name string) json.NewEncoder(w).Encode(sub) } -func (s *Server) deleteSubnet(w http.ResponseWriter, _ *http.Request, name string) { - cmd := dispatcher.DeleteSubnetCommand{Name: name} - if err := s.dispatcher.Prepare(cmd); err != nil { - 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") +func (s *Server) deleteSubnet(w http.ResponseWriter, r *http.Request, name string) { + s.dispatcher.Dispatch(dispatcher.DeleteSubnetCommand{Name: name}) w.WriteHeader(http.StatusAccepted) - json.NewEncoder(w).Encode(Subnet{Name: name, State: state}) + json.NewEncoder(w).Encode(Subnet{Name: name, State: "deleting"}) } diff --git a/internal/api/agent/subnets.go b/internal/api/agent/subnets.go index 9e99a45..18fb81f 100644 --- a/internal/api/agent/subnets.go +++ b/internal/api/agent/subnets.go @@ -74,42 +74,21 @@ 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"}) return } - cmd := dispatcher.CreateSubnetCommand{ + s.dispatcher.Dispatch(dispatcher.CreateSubnetCommand{ Name: req.Name, VPC: req.VPC, VxlanID: req.VxlanID, IfaceType: req.IfaceType, GatewayIP: req.GatewayIP, CIDR: req.CIDR, - } - if err := s.dispatcher.Prepare(cmd); err != nil { - w.WriteHeader(http.StatusConflict) - json.NewEncoder(w).Encode(ErrorResponse{Error: err.Error()}) - return - } - s.dispatcher.Dispatch(cmd) - entries, _ := kv.ListByPrefix(s.db, "subnet/"+req.Name+"/") - 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) + json.NewEncoder(w).Encode(Subnet{ + Name: req.Name, + State: "creating", + VPC: req.VPC, + VxlanID: req.VxlanID, + GatewayIP: req.GatewayIP, + CIDR: req.CIDR, + }) } diff --git a/internal/api/agent/vpc.go b/internal/api/agent/vpc.go index be0724e..fd1bfdb 100644 --- a/internal/api/agent/vpc.go +++ b/internal/api/agent/vpc.go @@ -40,14 +40,7 @@ func (s *Server) getVpc(w http.ResponseWriter, _ *http.Request, name string) { } func (s *Server) deleteVpc(w http.ResponseWriter, _ *http.Request, name string) { - cmd := dispatcher.DeleteVPCCommand{Name: name} - if err := s.dispatcher.Prepare(cmd); err != nil { - 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") + s.dispatcher.Dispatch(dispatcher.DeleteVPCCommand{Name: name}) w.WriteHeader(http.StatusAccepted) - json.NewEncoder(w).Encode(VPC{Name: name, State: state}) + json.NewEncoder(w).Encode(VPC{Name: name, State: "deleting"}) } diff --git a/internal/api/agent/vpcs.go b/internal/api/agent/vpcs.go index 3087456..735259c 100644 --- a/internal/api/agent/vpcs.go +++ b/internal/api/agent/vpcs.go @@ -62,14 +62,7 @@ func (s *Server) postVpc(w http.ResponseWriter, r *http.Request) { json.NewEncoder(w).Encode(ErrorResponse{Error: "name is required"}) return } - cmd := dispatcher.CreateVPCCommand{Name: req.Name} - if err := s.dispatcher.Prepare(cmd); err != nil { - 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") + s.dispatcher.Dispatch(dispatcher.CreateVPCCommand{Name: req.Name}) w.WriteHeader(http.StatusAccepted) - json.NewEncoder(w).Encode(VPC{Name: req.Name, State: state}) + json.NewEncoder(w).Encode(VPC{Name: req.Name, State: "creating"}) } diff --git a/internal/config/agent/struct.go b/internal/config/agent/struct.go index 8818541..1c9fc9f 100644 --- a/internal/config/agent/struct.go +++ b/internal/config/agent/struct.go @@ -20,10 +20,6 @@ type Config struct { Count int `mapstructure:"count"` BufferSize int `mapstructure:"buffer_size"` } `mapstructure:"worker"` - Dispatcher struct { - TimeoutSeconds int `mapstructure:"timeout_seconds"` - PollSeconds int `mapstructure:"poll_seconds"` - } `mapstructure:"dispatcher"` DefaultInterface string `mapstructure:"default_interface"` Interfaces map[string]string `mapstructure:"interfaces"` } @@ -40,8 +36,6 @@ func LoadConfig(path string) (*Config, error) { v.SetDefault("prometheus.port", 9090) v.SetDefault("worker.count", 4) 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.ReadInConfig() diff --git a/internal/dispatcher/agent/dispatcher.go b/internal/dispatcher/agent/dispatcher.go index 6b1ca67..8e94827 100644 --- a/internal/dispatcher/agent/dispatcher.go +++ b/internal/dispatcher/agent/dispatcher.go @@ -9,7 +9,6 @@ import ( ) type Command interface { - Prepare(db *badger.DB, cfg *configuration.Config) error Execute(db *badger.DB, cfg *configuration.Config) error } @@ -23,10 +22,6 @@ func New(queue *worker.Queue, db *badger.DB, cfg *configuration.Config) *Dispatc 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) { d.queue.Submit(func() { if err := cmd.Execute(d.db, d.cfg); err != nil { diff --git a/internal/dispatcher/agent/subnet_commands.go b/internal/dispatcher/agent/subnet_commands.go index 7bc214d..0d471c3 100644 --- a/internal/dispatcher/agent/subnet_commands.go +++ b/internal/dispatcher/agent/subnet_commands.go @@ -4,7 +4,6 @@ import ( "fmt" "os" "strconv" - "time" configuration "git.g3e.fr/syonad/two/internal/config/agent" "git.g3e.fr/syonad/two/internal/subnet" @@ -21,17 +20,7 @@ type CreateSubnetCommand struct { CIDR string } -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) - } +func (c CreateSubnetCommand) Execute(db *badger.DB, cfg *configuration.Config) error { localIface, ok := cfg.Interfaces[c.IfaceType] if !ok { localIface = cfg.DefaultInterface @@ -42,25 +31,6 @@ func (c CreateSubnetCommand) Prepare(db *badger.DB, cfg *configuration.Config) e 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+"/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) } @@ -68,14 +38,8 @@ type DeleteSubnetCommand struct { 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 { + kv.AddInDB(db, "subnet/"+c.Name+"/state", "deleting") if err := subnet.DeleteSubnet(db, c.Name); err != nil { fmt.Println(err) os.Exit(1) diff --git a/internal/dispatcher/agent/vpc_commands.go b/internal/dispatcher/agent/vpc_commands.go index c03dd77..b2687fc 100644 --- a/internal/dispatcher/agent/vpc_commands.go +++ b/internal/dispatcher/agent/vpc_commands.go @@ -1,10 +1,6 @@ package dispatcher import ( - "fmt" - "strings" - "time" - configuration "git.g3e.fr/syonad/two/internal/config/agent" "git.g3e.fr/syonad/two/internal/vpc" "git.g3e.fr/syonad/two/pkg/db/kv" @@ -15,14 +11,8 @@ type CreateVPCCommand struct { 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 { + kv.AddInDB(db, "vpc/"+c.Name+"/state", "creating") return vpc.CreateVPC(db, c.Name) } @@ -30,54 +20,8 @@ type DeleteVPCCommand struct { Name string } -func (c DeleteVPCCommand) Prepare(db *badger.DB, _ *configuration.Config) error { - 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): - } - } +func (c DeleteVPCCommand) Execute(db *badger.DB, _ *configuration.Config) error { + kv.AddInDB(db, "vpc/"+c.Name+"/state", "deleting") if err := vpc.DeleteVPC(db, c.Name); err != nil { return err }