Compare commits
4 commits
53ca41e635
...
b0c73b2952
| Author | SHA1 | Date | |
|---|---|---|---|
|
b0c73b2952 |
|||
|
75e02bd4f7 |
|||
|
24536ed891 |
|||
|
8637c33bd8 |
10 changed files with 219 additions and 36 deletions
|
|
@ -6,10 +6,11 @@ import (
|
||||||
"log"
|
"log"
|
||||||
|
|
||||||
agentapi "git.g3e.fr/syonad/two/internal/api/agent"
|
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"
|
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"
|
"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"
|
"github.com/prometheus/client_golang/prometheus"
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
@ -22,16 +23,19 @@ func main() {
|
||||||
log.Fatalf("failed to load config: %v", err)
|
log.Fatalf("failed to load config: %v", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
db := kv.InitDB(kv.Config{Path: cfg.Database.Path}, true)
|
db := kv.InitDB(kv.Config{Path: cfg.Database.Path}, false)
|
||||||
defer db.Close()
|
defer db.Close()
|
||||||
|
|
||||||
apiAddr := fmt.Sprintf("%s:%d", cfg.Api.Address, cfg.Api.Port)
|
q := worker.New(cfg.Worker.BufferSize)
|
||||||
promAddr := fmt.Sprintf("%s:%d", cfg.Prometheus.Address, cfg.Prometheus.Port)
|
q.Start(cfg.Worker.Count)
|
||||||
|
|
||||||
registry := prometheus.NewRegistry()
|
registry := prometheus.NewRegistry()
|
||||||
registry.MustRegister(agentmetrics.NewAgentCollector(db))
|
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)
|
go promserver.Start(promAddr, registry)
|
||||||
|
|
||||||
select {}
|
select {}
|
||||||
|
|
|
||||||
33
internal/api/agent/models.go
Normal file
33
internal/api/agent/models.go
Normal file
|
|
@ -0,0 +1,33 @@
|
||||||
|
package agentapi
|
||||||
|
|
||||||
|
type VPCCreateRequest struct {
|
||||||
|
Name string `json:"name"`
|
||||||
|
}
|
||||||
|
|
||||||
|
type VPC struct {
|
||||||
|
Name string `json:"name"`
|
||||||
|
State string `json:"state"`
|
||||||
|
}
|
||||||
|
|
||||||
|
type SubnetCreateRequest struct {
|
||||||
|
Name string `json:"name"`
|
||||||
|
VPC string `json:"vpc"`
|
||||||
|
VxlanID int `json:"vxlan_id"`
|
||||||
|
LocalIP string `json:"local_ip"`
|
||||||
|
GatewayIP string `json:"gateway_ip"`
|
||||||
|
CIDR string `json:"cidr"`
|
||||||
|
}
|
||||||
|
|
||||||
|
type Subnet struct {
|
||||||
|
Name string `json:"name"`
|
||||||
|
State string `json:"state"`
|
||||||
|
VPC string `json:"vpc"`
|
||||||
|
VxlanID int `json:"vxlan_id"`
|
||||||
|
LocalIP string `json:"local_ip"`
|
||||||
|
GatewayIP string `json:"gateway_ip"`
|
||||||
|
CIDR string `json:"cidr"`
|
||||||
|
}
|
||||||
|
|
||||||
|
type ErrorResponse struct {
|
||||||
|
Error string `json:"error"`
|
||||||
|
}
|
||||||
|
|
@ -3,14 +3,31 @@ package agentapi
|
||||||
import (
|
import (
|
||||||
"log"
|
"log"
|
||||||
"net/http"
|
"net/http"
|
||||||
|
|
||||||
|
"git.g3e.fr/syonad/two/pkg/worker"
|
||||||
)
|
)
|
||||||
|
|
||||||
func Start(address string) {
|
type Server struct {
|
||||||
mux := http.NewServeMux()
|
queue *worker.Queue
|
||||||
mux.HandleFunc("/vpcs", VpcsHandler)
|
}
|
||||||
mux.HandleFunc("/vpcs/", VpcByNameHandler)
|
|
||||||
mux.HandleFunc("/subnets", SubnetsHandler)
|
func New(queue *worker.Queue) *Server {
|
||||||
mux.HandleFunc("/subnets/", SubnetByNameHandler)
|
return &Server{queue: queue}
|
||||||
log.Printf("API server listening on %s", address)
|
}
|
||||||
log.Fatal(http.ListenAndServe(address, mux))
|
|
||||||
|
func (s *Server) Start(address string) {
|
||||||
|
mux := http.NewServeMux()
|
||||||
|
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, logMiddleware(mux)))
|
||||||
|
}
|
||||||
|
|
||||||
|
func logMiddleware(next http.Handler) http.Handler {
|
||||||
|
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||||
|
log.Printf("%s %s %s", r.RemoteAddr, r.Method, r.URL.Path)
|
||||||
|
next.ServeHTTP(w, r)
|
||||||
|
})
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -6,20 +6,36 @@ import (
|
||||||
"strings"
|
"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/")
|
name := strings.TrimPrefix(r.URL.Path, "/subnets/")
|
||||||
if name == "" {
|
if name == "" {
|
||||||
http.NotFound(w, r)
|
w.Header().Set("Content-Type", "application/json")
|
||||||
|
w.WriteHeader(http.StatusNotFound)
|
||||||
|
json.NewEncoder(w).Encode(ErrorResponse{Error: "resource not found"})
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
w.Header().Set("Content-Type", "application/json")
|
w.Header().Set("Content-Type", "application/json")
|
||||||
switch r.Method {
|
switch r.Method {
|
||||||
case http.MethodGet:
|
case http.MethodGet:
|
||||||
w.WriteHeader(http.StatusOK)
|
s.getSubnet(w, r, name)
|
||||||
json.NewEncoder(w).Encode(map[string]string{"name": name})
|
|
||||||
case http.MethodDelete:
|
case http.MethodDelete:
|
||||||
w.WriteHeader(http.StatusAccepted)
|
s.deleteSubnet(w, r, name)
|
||||||
default:
|
default:
|
||||||
http.Error(w, "method not allowed", http.StatusMethodNotAllowed)
|
http.Error(w, `{"error":"method not allowed"}`, http.StatusMethodNotAllowed)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (s *Server) getSubnet(w http.ResponseWriter, r *http.Request, name string) {
|
||||||
|
w.WriteHeader(http.StatusOK)
|
||||||
|
json.NewEncoder(w).Encode(Subnet{Name: name, State: "created"})
|
||||||
|
}
|
||||||
|
|
||||||
|
func (s *Server) deleteSubnet(w http.ResponseWriter, r *http.Request, name string) {
|
||||||
|
s.queue.Submit(func() {
|
||||||
|
destroySubnet(name)
|
||||||
|
})
|
||||||
|
w.WriteHeader(http.StatusAccepted)
|
||||||
|
json.NewEncoder(w).Encode(Subnet{Name: name, State: "deleting"})
|
||||||
|
}
|
||||||
|
|
||||||
|
func destroySubnet(name string) {}
|
||||||
|
|
|
||||||
|
|
@ -5,15 +5,48 @@ import (
|
||||||
"net/http"
|
"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")
|
w.Header().Set("Content-Type", "application/json")
|
||||||
switch r.Method {
|
switch r.Method {
|
||||||
case http.MethodGet:
|
case http.MethodGet:
|
||||||
w.WriteHeader(http.StatusOK)
|
s.listSubnets(w, r)
|
||||||
json.NewEncoder(w).Encode([]interface{}{})
|
|
||||||
case http.MethodPost:
|
case http.MethodPost:
|
||||||
w.WriteHeader(http.StatusAccepted)
|
s.postSubnet(w, r)
|
||||||
default:
|
default:
|
||||||
http.Error(w, "method not allowed", http.StatusMethodNotAllowed)
|
http.Error(w, `{"error":"method not allowed"}`, http.StatusMethodNotAllowed)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (s *Server) listSubnets(w http.ResponseWriter, r *http.Request) {
|
||||||
|
w.WriteHeader(http.StatusOK)
|
||||||
|
json.NewEncoder(w).Encode([]Subnet{})
|
||||||
|
}
|
||||||
|
|
||||||
|
func (s *Server) postSubnet(w http.ResponseWriter, r *http.Request) {
|
||||||
|
var req SubnetCreateRequest
|
||||||
|
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
|
||||||
|
w.WriteHeader(http.StatusBadRequest)
|
||||||
|
json.NewEncoder(w).Encode(ErrorResponse{Error: "invalid request body"})
|
||||||
|
return
|
||||||
|
}
|
||||||
|
if req.Name == "" || req.VPC == "" || req.LocalIP == "" || req.GatewayIP == "" || req.CIDR == "" {
|
||||||
|
w.WriteHeader(http.StatusBadRequest)
|
||||||
|
json.NewEncoder(w).Encode(ErrorResponse{Error: "name, vpc, local_ip, gateway_ip and cidr are required"})
|
||||||
|
return
|
||||||
|
}
|
||||||
|
s.queue.Submit(func() {
|
||||||
|
createSubnet(req)
|
||||||
|
})
|
||||||
|
w.WriteHeader(http.StatusAccepted)
|
||||||
|
json.NewEncoder(w).Encode(Subnet{
|
||||||
|
Name: req.Name,
|
||||||
|
State: "creating",
|
||||||
|
VPC: req.VPC,
|
||||||
|
VxlanID: req.VxlanID,
|
||||||
|
LocalIP: req.LocalIP,
|
||||||
|
GatewayIP: req.GatewayIP,
|
||||||
|
CIDR: req.CIDR,
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
func createSubnet(req SubnetCreateRequest) {}
|
||||||
|
|
|
||||||
|
|
@ -6,20 +6,36 @@ import (
|
||||||
"strings"
|
"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/")
|
name := strings.TrimPrefix(r.URL.Path, "/vpcs/")
|
||||||
if name == "" {
|
if name == "" {
|
||||||
http.NotFound(w, r)
|
w.Header().Set("Content-Type", "application/json")
|
||||||
|
w.WriteHeader(http.StatusNotFound)
|
||||||
|
json.NewEncoder(w).Encode(ErrorResponse{Error: "resource not found"})
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
w.Header().Set("Content-Type", "application/json")
|
w.Header().Set("Content-Type", "application/json")
|
||||||
switch r.Method {
|
switch r.Method {
|
||||||
case http.MethodGet:
|
case http.MethodGet:
|
||||||
w.WriteHeader(http.StatusOK)
|
s.getVpc(w, r, name)
|
||||||
json.NewEncoder(w).Encode(map[string]string{"name": name})
|
|
||||||
case http.MethodDelete:
|
case http.MethodDelete:
|
||||||
w.WriteHeader(http.StatusAccepted)
|
s.deleteVpc(w, r, name)
|
||||||
default:
|
default:
|
||||||
http.Error(w, "method not allowed", http.StatusMethodNotAllowed)
|
http.Error(w, `{"error":"method not allowed"}`, http.StatusMethodNotAllowed)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (s *Server) getVpc(w http.ResponseWriter, r *http.Request, name string) {
|
||||||
|
w.WriteHeader(http.StatusOK)
|
||||||
|
json.NewEncoder(w).Encode(VPC{Name: name, State: "created"})
|
||||||
|
}
|
||||||
|
|
||||||
|
func (s *Server) deleteVpc(w http.ResponseWriter, r *http.Request, name string) {
|
||||||
|
s.queue.Submit(func() {
|
||||||
|
destroyVpc(name)
|
||||||
|
})
|
||||||
|
w.WriteHeader(http.StatusAccepted)
|
||||||
|
json.NewEncoder(w).Encode(VPC{Name: name, State: "deleting"})
|
||||||
|
}
|
||||||
|
|
||||||
|
func destroyVpc(name string) {}
|
||||||
|
|
|
||||||
|
|
@ -5,15 +5,40 @@ import (
|
||||||
"net/http"
|
"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")
|
w.Header().Set("Content-Type", "application/json")
|
||||||
switch r.Method {
|
switch r.Method {
|
||||||
case http.MethodGet:
|
case http.MethodGet:
|
||||||
w.WriteHeader(http.StatusOK)
|
s.listVpcs(w, r)
|
||||||
json.NewEncoder(w).Encode([]interface{}{})
|
|
||||||
case http.MethodPost:
|
case http.MethodPost:
|
||||||
w.WriteHeader(http.StatusAccepted)
|
s.postVpc(w, r)
|
||||||
default:
|
default:
|
||||||
http.Error(w, "method not allowed", http.StatusMethodNotAllowed)
|
http.Error(w, `{"error":"method not allowed"}`, http.StatusMethodNotAllowed)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (s *Server) listVpcs(w http.ResponseWriter, r *http.Request) {
|
||||||
|
w.WriteHeader(http.StatusOK)
|
||||||
|
json.NewEncoder(w).Encode([]VPC{})
|
||||||
|
}
|
||||||
|
|
||||||
|
func (s *Server) postVpc(w http.ResponseWriter, r *http.Request) {
|
||||||
|
var req VPCCreateRequest
|
||||||
|
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
|
||||||
|
w.WriteHeader(http.StatusBadRequest)
|
||||||
|
json.NewEncoder(w).Encode(ErrorResponse{Error: "invalid request body"})
|
||||||
|
return
|
||||||
|
}
|
||||||
|
if req.Name == "" {
|
||||||
|
w.WriteHeader(http.StatusBadRequest)
|
||||||
|
json.NewEncoder(w).Encode(ErrorResponse{Error: "name is required"})
|
||||||
|
return
|
||||||
|
}
|
||||||
|
s.queue.Submit(func() {
|
||||||
|
createVpc(req.Name)
|
||||||
|
})
|
||||||
|
w.WriteHeader(http.StatusAccepted)
|
||||||
|
json.NewEncoder(w).Encode(VPC{Name: req.Name, State: "creating"})
|
||||||
|
}
|
||||||
|
|
||||||
|
func createVpc(name string) {}
|
||||||
|
|
|
||||||
|
|
@ -16,6 +16,10 @@ type Config struct {
|
||||||
Address string `mapstructure:"address"`
|
Address string `mapstructure:"address"`
|
||||||
Port int `mapstructure:"port"`
|
Port int `mapstructure:"port"`
|
||||||
} `mapstructure:"prometheus"`
|
} `mapstructure:"prometheus"`
|
||||||
|
Worker struct {
|
||||||
|
Count int `mapstructure:"count"`
|
||||||
|
BufferSize int `mapstructure:"buffer_size"`
|
||||||
|
} `mapstructure:"worker"`
|
||||||
}
|
}
|
||||||
|
|
||||||
func LoadConfig(path string) (*Config, error) {
|
func LoadConfig(path string) (*Config, error) {
|
||||||
|
|
@ -28,6 +32,8 @@ func LoadConfig(path string) (*Config, error) {
|
||||||
v.SetDefault("api.port", 8080)
|
v.SetDefault("api.port", 8080)
|
||||||
v.SetDefault("prometheus.address", "")
|
v.SetDefault("prometheus.address", "")
|
||||||
v.SetDefault("prometheus.port", 9090)
|
v.SetDefault("prometheus.port", 9090)
|
||||||
|
v.SetDefault("worker.count", 4)
|
||||||
|
v.SetDefault("worker.buffer_size", 100)
|
||||||
|
|
||||||
v.ReadInConfig()
|
v.ReadInConfig()
|
||||||
|
|
||||||
|
|
|
||||||
33
pkg/worker/queue.go
Normal file
33
pkg/worker/queue.go
Normal file
|
|
@ -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)
|
||||||
|
}
|
||||||
|
}
|
||||||
Loading…
Add table
Add a link
Reference in a new issue