f-21: dispatch: add check and verif before and during action

Signed-off-by: GnomeZworc <nicolas.boufidjeline@g3e.fr>
This commit is contained in:
GnomeZworc 2026-04-23 23:29:38 +02:00
commit 9ee361ae28
Signed by: nicolas.boufideline
GPG key ID: 4406BBBF8845D632
3 changed files with 66 additions and 2 deletions

View file

@ -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()

View file

@ -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"
@ -44,7 +45,22 @@ func (c CreateSubnetCommand) Prepare(db *badger.DB, cfg *configuration.Config) e
return nil return nil
} }
func (c CreateSubnetCommand) Execute(db *badger.DB, _ *configuration.Config) error { 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)
} }

View file

@ -2,6 +2,8 @@ package dispatcher
import ( import (
"fmt" "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"
@ -32,10 +34,50 @@ func (c DeleteVPCCommand) 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 not found", c.Name) 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") return kv.AddInDB(db, "vpc/"+c.Name+"/state", "deleting")
} }
func (c DeleteVPCCommand) Execute(db *badger.DB, _ *configuration.Config) error { 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
} }