From 63a288f69eafc0f41ca6550fa06543cc816af801 Mon Sep 17 00:00:00 2001 From: GnomeZworc Date: Sun, 19 Apr 2026 00:06:50 +0200 Subject: [PATCH] f-21: refactor: add dispatcher layer for MQTT migration Introduce internal/dispatcher package with a Command interface and typed commands (CreateVPC, DeleteVPC, CreateSubnet, DeleteSubnet). The API handlers now call dispatcher.Dispatch() instead of enqueuing closures directly, decoupling transport (HTTP today, MQTT tomorrow) from execution. Signed-off-by: GnomeZworc --- cmd/agent/main.go | 4 ++- internal/api/agent/server.go | 10 ++++---- internal/api/agent/subnet.go | 7 +++--- internal/api/agent/subnets.go | 12 ++++++--- internal/api/agent/vpc.go | 19 ++------------ internal/api/agent/vpcs.go | 11 ++------ internal/dispatcher/dispatcher.go | 29 +++++++++++++++++++++ internal/dispatcher/subnet_commands.go | 26 +++++++++++++++++++ internal/dispatcher/vpc_commands.go | 35 ++++++++++++++++++++++++++ 9 files changed, 114 insertions(+), 39 deletions(-) create mode 100644 internal/dispatcher/dispatcher.go create mode 100644 internal/dispatcher/subnet_commands.go create mode 100644 internal/dispatcher/vpc_commands.go diff --git a/cmd/agent/main.go b/cmd/agent/main.go index 2f5e084..b61bb6f 100644 --- a/cmd/agent/main.go +++ b/cmd/agent/main.go @@ -7,6 +7,7 @@ import ( agentapi "git.g3e.fr/syonad/two/internal/api/agent" configuration "git.g3e.fr/syonad/two/internal/config/agent" + "git.g3e.fr/syonad/two/internal/dispatcher" agentmetrics "git.g3e.fr/syonad/two/internal/prometheus/agent" "git.g3e.fr/syonad/two/pkg/db/kv" promserver "git.g3e.fr/syonad/two/pkg/prometheus" @@ -35,7 +36,8 @@ func main() { apiAddr := fmt.Sprintf("%s:%d", cfg.Api.Address, cfg.Api.Port) promAddr := fmt.Sprintf("%s:%d", cfg.Prometheus.Address, cfg.Prometheus.Port) - go agentapi.New(q, db).Start(apiAddr) + d := dispatcher.New(q, db) + go agentapi.New(d, db).Start(apiAddr) go promserver.Start(promAddr, registry) select {} diff --git a/internal/api/agent/server.go b/internal/api/agent/server.go index 7e4eb21..7f3247c 100644 --- a/internal/api/agent/server.go +++ b/internal/api/agent/server.go @@ -4,17 +4,17 @@ import ( "log" "net/http" - "git.g3e.fr/syonad/two/pkg/worker" + "git.g3e.fr/syonad/two/internal/dispatcher" "github.com/dgraph-io/badger/v4" ) type Server struct { - queue *worker.Queue - db *badger.DB + dispatcher *dispatcher.Dispatcher + db *badger.DB } -func New(queue *worker.Queue, db *badger.DB) *Server { - return &Server{queue: queue, db: db} +func New(d *dispatcher.Dispatcher, db *badger.DB) *Server { + return &Server{dispatcher: d, db: db} } func (s *Server) Start(address string) { diff --git a/internal/api/agent/subnet.go b/internal/api/agent/subnet.go index 289667f..b20c65a 100644 --- a/internal/api/agent/subnet.go +++ b/internal/api/agent/subnet.go @@ -4,6 +4,8 @@ import ( "encoding/json" "net/http" "strings" + + "git.g3e.fr/syonad/two/internal/dispatcher" ) func (s *Server) SubnetByNameHandler(w http.ResponseWriter, r *http.Request) { @@ -31,11 +33,8 @@ func (s *Server) getSubnet(w http.ResponseWriter, r *http.Request, name string) } func (s *Server) deleteSubnet(w http.ResponseWriter, r *http.Request, name string) { - s.queue.Submit(func() { - destroySubnet(name) - }) + s.dispatcher.Dispatch(dispatcher.DeleteSubnetCommand{Name: name}) w.WriteHeader(http.StatusAccepted) json.NewEncoder(w).Encode(Subnet{Name: name, State: "deleting"}) } -func destroySubnet(name string) {} diff --git a/internal/api/agent/subnets.go b/internal/api/agent/subnets.go index 5e0ebdc..b251c3c 100644 --- a/internal/api/agent/subnets.go +++ b/internal/api/agent/subnets.go @@ -3,6 +3,8 @@ package agentapi import ( "encoding/json" "net/http" + + "git.g3e.fr/syonad/two/internal/dispatcher" ) func (s *Server) SubnetsHandler(w http.ResponseWriter, r *http.Request) { @@ -34,8 +36,13 @@ func (s *Server) postSubnet(w http.ResponseWriter, r *http.Request) { json.NewEncoder(w).Encode(ErrorResponse{Error: "name, vpc, local_ip, gateway_ip and cidr are required"}) return } - s.queue.Submit(func() { - createSubnet(req) + s.dispatcher.Dispatch(dispatcher.CreateSubnetCommand{ + Name: req.Name, + VPC: req.VPC, + VxlanID: req.VxlanID, + LocalIP: req.LocalIP, + GatewayIP: req.GatewayIP, + CIDR: req.CIDR, }) w.WriteHeader(http.StatusAccepted) json.NewEncoder(w).Encode(Subnet{ @@ -49,4 +56,3 @@ func (s *Server) postSubnet(w http.ResponseWriter, r *http.Request) { }) } -func createSubnet(req SubnetCreateRequest) {} diff --git a/internal/api/agent/vpc.go b/internal/api/agent/vpc.go index 3396d23..60d614e 100644 --- a/internal/api/agent/vpc.go +++ b/internal/api/agent/vpc.go @@ -2,13 +2,10 @@ package agentapi import ( "encoding/json" - "fmt" "net/http" - "os" "strings" - "git.g3e.fr/syonad/two/internal/vpc" - "git.g3e.fr/syonad/two/pkg/db/kv" + "git.g3e.fr/syonad/two/internal/dispatcher" ) func (s *Server) VpcByNameHandler(w http.ResponseWriter, r *http.Request) { @@ -36,20 +33,8 @@ func (s *Server) getVpc(w http.ResponseWriter, _ *http.Request, name string) { } func (s *Server) deleteVpc(w http.ResponseWriter, _ *http.Request, name string) { - s.queue.Submit(func() { - kv.AddInDB(s.db, "vpc/"+name+"/state", "deleting") - if err := vpc.DeleteVPC(s.db, name); err != nil { - fmt.Println(err) - } - if state, err := kv.GetFromDB(s.db, "vpc/"+name+"/state"); err != nil { - fmt.Println(err) - os.Exit(1) - } else if state == "deleted" { - kv.DeleteInDB(s.db, "vpc/"+name) - } - }) + s.dispatcher.Dispatch(dispatcher.DeleteVPCCommand{Name: name}) w.WriteHeader(http.StatusAccepted) json.NewEncoder(w).Encode(VPC{Name: name, State: "deleting"}) } -func destroyVpc(name string) {} diff --git a/internal/api/agent/vpcs.go b/internal/api/agent/vpcs.go index 2dddd42..eb633e3 100644 --- a/internal/api/agent/vpcs.go +++ b/internal/api/agent/vpcs.go @@ -2,11 +2,9 @@ package agentapi import ( "encoding/json" - "fmt" "net/http" - "git.g3e.fr/syonad/two/internal/vpc" - "git.g3e.fr/syonad/two/pkg/db/kv" + "git.g3e.fr/syonad/two/internal/dispatcher" ) func (s *Server) VpcsHandler(w http.ResponseWriter, r *http.Request) { @@ -38,12 +36,7 @@ func (s *Server) postVpc(w http.ResponseWriter, r *http.Request) { json.NewEncoder(w).Encode(ErrorResponse{Error: "name is required"}) return } - s.queue.Submit(func() { - kv.AddInDB(s.db, "vpc/"+req.Name+"/state", "creating") - if err := vpc.CreateVPC(s.db, req.Name); err != nil { - fmt.Println(err) - } - }) + s.dispatcher.Dispatch(dispatcher.CreateVPCCommand{Name: req.Name}) w.WriteHeader(http.StatusAccepted) json.NewEncoder(w).Encode(VPC{Name: req.Name, State: "creating"}) } diff --git a/internal/dispatcher/dispatcher.go b/internal/dispatcher/dispatcher.go new file mode 100644 index 0000000..937adf2 --- /dev/null +++ b/internal/dispatcher/dispatcher.go @@ -0,0 +1,29 @@ +package dispatcher + +import ( + "log" + + "git.g3e.fr/syonad/two/pkg/worker" + "github.com/dgraph-io/badger/v4" +) + +type Command interface { + Execute(db *badger.DB) error +} + +type Dispatcher struct { + queue *worker.Queue + db *badger.DB +} + +func New(queue *worker.Queue, db *badger.DB) *Dispatcher { + return &Dispatcher{queue: queue, db: db} +} + +func (d *Dispatcher) Dispatch(cmd Command) { + d.queue.Submit(func() { + if err := cmd.Execute(d.db); err != nil { + log.Printf("command error (%T): %v", cmd, err) + } + }) +} diff --git a/internal/dispatcher/subnet_commands.go b/internal/dispatcher/subnet_commands.go new file mode 100644 index 0000000..d8ef8e3 --- /dev/null +++ b/internal/dispatcher/subnet_commands.go @@ -0,0 +1,26 @@ +package dispatcher + +import "github.com/dgraph-io/badger/v4" + +type CreateSubnetCommand struct { + Name string + VPC string + VxlanID int + LocalIP string + GatewayIP string + CIDR string +} + +func (c CreateSubnetCommand) Execute(db *badger.DB) error { + // TODO: brancher internal/subnet/create.go + return nil +} + +type DeleteSubnetCommand struct { + Name string +} + +func (c DeleteSubnetCommand) Execute(db *badger.DB) error { + // TODO: brancher internal/subnet/delete.go + return nil +} diff --git a/internal/dispatcher/vpc_commands.go b/internal/dispatcher/vpc_commands.go new file mode 100644 index 0000000..195b014 --- /dev/null +++ b/internal/dispatcher/vpc_commands.go @@ -0,0 +1,35 @@ +package dispatcher + +import ( + "git.g3e.fr/syonad/two/internal/vpc" + "git.g3e.fr/syonad/two/pkg/db/kv" + "github.com/dgraph-io/badger/v4" +) + +type CreateVPCCommand struct { + Name string +} + +func (c CreateVPCCommand) Execute(db *badger.DB) error { + kv.AddInDB(db, "vpc/"+c.Name+"/state", "creating") + return vpc.CreateVPC(db, c.Name) +} + +type DeleteVPCCommand struct { + Name string +} + +func (c DeleteVPCCommand) Execute(db *badger.DB) error { + kv.AddInDB(db, "vpc/"+c.Name+"/state", "deleting") + if err := vpc.DeleteVPC(db, c.Name); err != nil { + return err + } + state, err := kv.GetFromDB(db, "vpc/"+c.Name+"/state") + if err != nil { + return err + } + if state == "deleted" { + kv.DeleteInDB(db, "vpc/"+c.Name) + } + return nil +}