diff --git a/agent b/agent new file mode 100755 index 0000000..e6b9fa4 Binary files /dev/null and b/agent differ diff --git a/cmd/agent/main.go b/cmd/agent/main.go index ce3df4e..56011bf 100644 --- a/cmd/agent/main.go +++ b/cmd/agent/main.go @@ -6,10 +6,11 @@ import ( "log" agentapi "git.g3e.fr/syonad/two/internal/api/agent" - agentmetrics "git.g3e.fr/syonad/two/internal/prometheus/agent" configuration "git.g3e.fr/syonad/two/internal/config/agent" - promserver "git.g3e.fr/syonad/two/pkg/prometheus" + 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" + "git.g3e.fr/syonad/two/pkg/worker" "github.com/prometheus/client_golang/prometheus" ) @@ -25,13 +26,16 @@ func main() { db := kv.InitDB(kv.Config{Path: cfg.Database.Path}, true) defer db.Close() - apiAddr := fmt.Sprintf("%s:%d", cfg.Api.Address, cfg.Api.Port) - promAddr := fmt.Sprintf("%s:%d", cfg.Prometheus.Address, cfg.Prometheus.Port) + q := worker.New(cfg.Worker.BufferSize) + q.Start(cfg.Worker.Count) registry := prometheus.NewRegistry() registry.MustRegister(agentmetrics.NewAgentCollector(db)) - go agentapi.Start(apiAddr) + 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).Start(apiAddr) go promserver.Start(promAddr, registry) select {} diff --git a/internal/api/agent/server.go b/internal/api/agent/server.go index 05653a4..aebe2c3 100644 --- a/internal/api/agent/server.go +++ b/internal/api/agent/server.go @@ -3,14 +3,24 @@ package agentapi import ( "log" "net/http" + + "git.g3e.fr/syonad/two/pkg/worker" ) -func Start(address string) { +type Server struct { + queue *worker.Queue +} + +func New(queue *worker.Queue) *Server { + return &Server{queue: queue} +} + +func (s *Server) Start(address string) { mux := http.NewServeMux() - mux.HandleFunc("/vpcs", VpcsHandler) - mux.HandleFunc("/vpcs/", VpcByNameHandler) - mux.HandleFunc("/subnets", SubnetsHandler) - mux.HandleFunc("/subnets/", SubnetByNameHandler) + mux.HandleFunc("/vpcs", s.VpcsHandler) + mux.HandleFunc("/vpcs/", s.VpcByNameHandler) + mux.HandleFunc("/subnets", s.SubnetsHandler) + mux.HandleFunc("/subnets/", s.SubnetByNameHandler) log.Printf("API server listening on %s", address) log.Fatal(http.ListenAndServe(address, mux)) } diff --git a/internal/api/agent/subnet.go b/internal/api/agent/subnet.go index e46e4f4..52b27d0 100644 --- a/internal/api/agent/subnet.go +++ b/internal/api/agent/subnet.go @@ -6,7 +6,7 @@ import ( "strings" ) -func SubnetByNameHandler(w http.ResponseWriter, r *http.Request) { +func (s *Server) SubnetByNameHandler(w http.ResponseWriter, r *http.Request) { name := strings.TrimPrefix(r.URL.Path, "/subnets/") if name == "" { http.NotFound(w, r) @@ -18,8 +18,13 @@ func SubnetByNameHandler(w http.ResponseWriter, r *http.Request) { w.WriteHeader(http.StatusOK) json.NewEncoder(w).Encode(map[string]string{"name": name}) case http.MethodDelete: + s.queue.Submit(func() { + deleteSubnet(name) + }) w.WriteHeader(http.StatusAccepted) default: http.Error(w, "method not allowed", http.StatusMethodNotAllowed) } } + +func deleteSubnet(name string) {} diff --git a/internal/api/agent/subnets.go b/internal/api/agent/subnets.go index 53dce2e..b3d587b 100644 --- a/internal/api/agent/subnets.go +++ b/internal/api/agent/subnets.go @@ -5,15 +5,20 @@ import ( "net/http" ) -func SubnetsHandler(w http.ResponseWriter, r *http.Request) { +func (s *Server) SubnetsHandler(w http.ResponseWriter, r *http.Request) { w.Header().Set("Content-Type", "application/json") switch r.Method { case http.MethodGet: w.WriteHeader(http.StatusOK) json.NewEncoder(w).Encode([]interface{}{}) case http.MethodPost: + s.queue.Submit(func() { + createSubnet() + }) w.WriteHeader(http.StatusAccepted) default: http.Error(w, "method not allowed", http.StatusMethodNotAllowed) } } + +func createSubnet() {} diff --git a/internal/api/agent/vpc.go b/internal/api/agent/vpc.go index aecfaa1..73cecd9 100644 --- a/internal/api/agent/vpc.go +++ b/internal/api/agent/vpc.go @@ -6,7 +6,7 @@ import ( "strings" ) -func VpcByNameHandler(w http.ResponseWriter, r *http.Request) { +func (s *Server) VpcByNameHandler(w http.ResponseWriter, r *http.Request) { name := strings.TrimPrefix(r.URL.Path, "/vpcs/") if name == "" { http.NotFound(w, r) @@ -18,8 +18,13 @@ func VpcByNameHandler(w http.ResponseWriter, r *http.Request) { w.WriteHeader(http.StatusOK) json.NewEncoder(w).Encode(map[string]string{"name": name}) case http.MethodDelete: + s.queue.Submit(func() { + deleteVpc(name) + }) w.WriteHeader(http.StatusAccepted) default: http.Error(w, "method not allowed", http.StatusMethodNotAllowed) } } + +func deleteVpc(name string) {} diff --git a/internal/api/agent/vpcs.go b/internal/api/agent/vpcs.go index c4e9b27..57386b6 100644 --- a/internal/api/agent/vpcs.go +++ b/internal/api/agent/vpcs.go @@ -5,15 +5,20 @@ import ( "net/http" ) -func VpcsHandler(w http.ResponseWriter, r *http.Request) { +func (s *Server) VpcsHandler(w http.ResponseWriter, r *http.Request) { w.Header().Set("Content-Type", "application/json") switch r.Method { case http.MethodGet: w.WriteHeader(http.StatusOK) json.NewEncoder(w).Encode([]interface{}{}) case http.MethodPost: + s.queue.Submit(func() { + createVpc() + }) w.WriteHeader(http.StatusAccepted) default: http.Error(w, "method not allowed", http.StatusMethodNotAllowed) } } + +func createVpc() {} diff --git a/internal/config/agent/struct.go b/internal/config/agent/struct.go index 0cca127..24ad0f5 100644 --- a/internal/config/agent/struct.go +++ b/internal/config/agent/struct.go @@ -16,6 +16,10 @@ type Config struct { Address string `mapstructure:"address"` Port int `mapstructure:"port"` } `mapstructure:"prometheus"` + Worker struct { + Count int `mapstructure:"count"` + BufferSize int `mapstructure:"buffer_size"` + } `mapstructure:"worker"` } func LoadConfig(path string) (*Config, error) { @@ -28,6 +32,8 @@ func LoadConfig(path string) (*Config, error) { v.SetDefault("api.port", 8080) v.SetDefault("prometheus.address", "") v.SetDefault("prometheus.port", 9090) + v.SetDefault("worker.count", 4) + v.SetDefault("worker.buffer_size", 100) v.ReadInConfig() diff --git a/pkg/worker/queue.go b/pkg/worker/queue.go new file mode 100644 index 0000000..109c726 --- /dev/null +++ b/pkg/worker/queue.go @@ -0,0 +1,33 @@ +package worker + +import "log" + +// Task is a function to be executed asynchronously by a worker. +type Task func() + +// Queue is a FIFO channel-backed task queue consumed by worker goroutines. +type Queue struct { + tasks chan Task +} + +// New creates a Queue with the given channel buffer size. +func New(bufferSize int) *Queue { + return &Queue{tasks: make(chan Task, bufferSize)} +} + +// Submit enqueues a task. Blocks if the queue is full. +func (q *Queue) Submit(t Task) { + q.tasks <- t +} + +// Start launches n worker goroutines that consume and execute tasks. +func (q *Queue) Start(n int) { + log.Printf("worker: starting %d workers", n) + for i := range n { + go func(id int) { + for task := range q.tasks { + task() + } + }(i) + } +}