Compare commits
2 commits
648a64782c
...
921a5ca96e
| Author | SHA1 | Date | |
|---|---|---|---|
|
921a5ca96e |
|||
|
fe6792be0d |
7 changed files with 122 additions and 29 deletions
|
|
@ -3,13 +3,14 @@ package main
|
||||||
import (
|
import (
|
||||||
"flag"
|
"flag"
|
||||||
"fmt"
|
"fmt"
|
||||||
"log"
|
"log/slog"
|
||||||
|
|
||||||
agentapi "git.g3e.fr/syonad/two/internal/api/agent"
|
agentapi "git.g3e.fr/syonad/two/internal/api/agent"
|
||||||
configuration "git.g3e.fr/syonad/two/internal/config/agent"
|
configuration "git.g3e.fr/syonad/two/internal/config/agent"
|
||||||
dispatcher "git.g3e.fr/syonad/two/internal/dispatcher/agent"
|
dispatcher "git.g3e.fr/syonad/two/internal/dispatcher/agent"
|
||||||
agentmetrics "git.g3e.fr/syonad/two/internal/prometheus/agent"
|
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"
|
||||||
|
"git.g3e.fr/syonad/two/pkg/logger"
|
||||||
promserver "git.g3e.fr/syonad/two/pkg/prometheus"
|
promserver "git.g3e.fr/syonad/two/pkg/prometheus"
|
||||||
"git.g3e.fr/syonad/two/pkg/worker"
|
"git.g3e.fr/syonad/two/pkg/worker"
|
||||||
"github.com/prometheus/client_golang/prometheus"
|
"github.com/prometheus/client_golang/prometheus"
|
||||||
|
|
@ -21,9 +22,12 @@ func main() {
|
||||||
|
|
||||||
cfg, err := configuration.LoadConfig(*confFile)
|
cfg, err := configuration.LoadConfig(*confFile)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
log.Fatalf("failed to load config: %v", err)
|
slog.Error("failed to load config", "error", err)
|
||||||
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
|
log := logger.New(cfg.Logger.Level, cfg.Logger.Debug)
|
||||||
|
|
||||||
db := kv.InitDB(kv.Config{Path: cfg.Database.Path}, false)
|
db := kv.InitDB(kv.Config{Path: cfg.Database.Path}, false)
|
||||||
defer db.Close()
|
defer db.Close()
|
||||||
|
|
||||||
|
|
@ -36,8 +40,16 @@ func main() {
|
||||||
apiAddr := fmt.Sprintf("%s:%d", cfg.Api.Address, cfg.Api.Port)
|
apiAddr := fmt.Sprintf("%s:%d", cfg.Api.Address, cfg.Api.Port)
|
||||||
promAddr := fmt.Sprintf("%s:%d", cfg.Prometheus.Address, cfg.Prometheus.Port)
|
promAddr := fmt.Sprintf("%s:%d", cfg.Prometheus.Address, cfg.Prometheus.Port)
|
||||||
|
|
||||||
d := dispatcher.New(q, db, cfg)
|
log.Info("starting agent",
|
||||||
go agentapi.New(d, db).Start(apiAddr)
|
"api", apiAddr,
|
||||||
|
"prometheus", promAddr,
|
||||||
|
"workers", cfg.Worker.Count,
|
||||||
|
"log_level", cfg.Logger.Level,
|
||||||
|
"debug", cfg.Logger.Debug,
|
||||||
|
)
|
||||||
|
|
||||||
|
d := dispatcher.New(q, db, cfg, log.With(slog.String("component", "dispatcher")))
|
||||||
|
go agentapi.New(d, db, log.With(slog.String("component", "api"))).Start(apiAddr)
|
||||||
go promserver.Start(promAddr, registry)
|
go promserver.Start(promAddr, registry)
|
||||||
|
|
||||||
select {}
|
select {}
|
||||||
|
|
|
||||||
|
|
@ -34,3 +34,10 @@ interfaces:
|
||||||
vms: br-000000
|
vms: br-000000
|
||||||
internet: br-000000
|
internet: br-000000
|
||||||
admin: br-000000
|
admin: br-000000
|
||||||
|
|
||||||
|
# Logging configuration
|
||||||
|
logger:
|
||||||
|
# Log level: debug, info, warn, error (default: info)
|
||||||
|
level: info
|
||||||
|
# Force debug level regardless of level setting (default: false)
|
||||||
|
debug: false
|
||||||
|
|
|
||||||
|
|
@ -1,8 +1,11 @@
|
||||||
package agentapi
|
package agentapi
|
||||||
|
|
||||||
import (
|
import (
|
||||||
"log"
|
"crypto/rand"
|
||||||
|
"encoding/hex"
|
||||||
|
"log/slog"
|
||||||
"net/http"
|
"net/http"
|
||||||
|
"time"
|
||||||
|
|
||||||
dispatcher "git.g3e.fr/syonad/two/internal/dispatcher/agent"
|
dispatcher "git.g3e.fr/syonad/two/internal/dispatcher/agent"
|
||||||
"github.com/dgraph-io/badger/v4"
|
"github.com/dgraph-io/badger/v4"
|
||||||
|
|
@ -11,10 +14,11 @@ import (
|
||||||
type Server struct {
|
type Server struct {
|
||||||
dispatcher *dispatcher.Dispatcher
|
dispatcher *dispatcher.Dispatcher
|
||||||
db *badger.DB
|
db *badger.DB
|
||||||
|
logger *slog.Logger
|
||||||
}
|
}
|
||||||
|
|
||||||
func New(d *dispatcher.Dispatcher, db *badger.DB) *Server {
|
func New(d *dispatcher.Dispatcher, db *badger.DB, logger *slog.Logger) *Server {
|
||||||
return &Server{dispatcher: d, db: db}
|
return &Server{dispatcher: d, db: db, logger: logger}
|
||||||
}
|
}
|
||||||
|
|
||||||
func (s *Server) Start(address string) {
|
func (s *Server) Start(address string) {
|
||||||
|
|
@ -23,13 +27,39 @@ func (s *Server) Start(address string) {
|
||||||
mux.HandleFunc("/vpcs/", s.VpcByNameHandler)
|
mux.HandleFunc("/vpcs/", s.VpcByNameHandler)
|
||||||
mux.HandleFunc("/subnets", s.SubnetsHandler)
|
mux.HandleFunc("/subnets", s.SubnetsHandler)
|
||||||
mux.HandleFunc("/subnets/", s.SubnetByNameHandler)
|
mux.HandleFunc("/subnets/", s.SubnetByNameHandler)
|
||||||
log.Printf("API server listening on %s", address)
|
s.logger.Info("API server listening", "address", address)
|
||||||
log.Fatal(http.ListenAndServe(address, logMiddleware(mux)))
|
if err := http.ListenAndServe(address, s.logMiddleware(mux)); err != nil {
|
||||||
|
s.logger.Error("API server stopped", "error", err)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
func logMiddleware(next http.Handler) http.Handler {
|
type statusWriter struct {
|
||||||
|
http.ResponseWriter
|
||||||
|
status int
|
||||||
|
}
|
||||||
|
|
||||||
|
func (sw *statusWriter) WriteHeader(code int) {
|
||||||
|
sw.status = code
|
||||||
|
sw.ResponseWriter.WriteHeader(code)
|
||||||
|
}
|
||||||
|
|
||||||
|
func (s *Server) logMiddleware(next http.Handler) http.Handler {
|
||||||
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||||
log.Printf("%s %s %s", r.RemoteAddr, r.Method, r.URL.Path)
|
var b [4]byte
|
||||||
next.ServeHTTP(w, r)
|
rand.Read(b[:])
|
||||||
|
reqID := hex.EncodeToString(b[:])
|
||||||
|
|
||||||
|
sw := &statusWriter{ResponseWriter: w, status: http.StatusOK}
|
||||||
|
start := time.Now()
|
||||||
|
next.ServeHTTP(sw, r)
|
||||||
|
|
||||||
|
s.logger.Info("request",
|
||||||
|
"request_id", reqID,
|
||||||
|
"method", r.Method,
|
||||||
|
"path", r.URL.Path,
|
||||||
|
"status", sw.status,
|
||||||
|
"duration_ms", time.Since(start).Milliseconds(),
|
||||||
|
"remote", r.RemoteAddr,
|
||||||
|
)
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -24,6 +24,10 @@ type Config struct {
|
||||||
TimeoutSeconds int `mapstructure:"timeout_seconds"`
|
TimeoutSeconds int `mapstructure:"timeout_seconds"`
|
||||||
PollSeconds int `mapstructure:"poll_seconds"`
|
PollSeconds int `mapstructure:"poll_seconds"`
|
||||||
} `mapstructure:"dispatcher"`
|
} `mapstructure:"dispatcher"`
|
||||||
|
Logger struct {
|
||||||
|
Level string `mapstructure:"level"`
|
||||||
|
Debug bool `mapstructure:"debug"`
|
||||||
|
} `mapstructure:"logger"`
|
||||||
DefaultInterface string `mapstructure:"default_interface"`
|
DefaultInterface string `mapstructure:"default_interface"`
|
||||||
Interfaces map[string]string `mapstructure:"interfaces"`
|
Interfaces map[string]string `mapstructure:"interfaces"`
|
||||||
}
|
}
|
||||||
|
|
@ -43,6 +47,8 @@ func LoadConfig(path string) (*Config, error) {
|
||||||
v.SetDefault("dispatcher.timeout_seconds", 300)
|
v.SetDefault("dispatcher.timeout_seconds", 300)
|
||||||
v.SetDefault("dispatcher.poll_seconds", 2)
|
v.SetDefault("dispatcher.poll_seconds", 2)
|
||||||
v.SetDefault("default_interface", "br-000000")
|
v.SetDefault("default_interface", "br-000000")
|
||||||
|
v.SetDefault("logger.level", "info")
|
||||||
|
v.SetDefault("logger.debug", false)
|
||||||
|
|
||||||
v.ReadInConfig()
|
v.ReadInConfig()
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -1,7 +1,9 @@
|
||||||
package dispatcher
|
package dispatcher
|
||||||
|
|
||||||
import (
|
import (
|
||||||
"log"
|
"fmt"
|
||||||
|
"log/slog"
|
||||||
|
"time"
|
||||||
|
|
||||||
configuration "git.g3e.fr/syonad/two/internal/config/agent"
|
configuration "git.g3e.fr/syonad/two/internal/config/agent"
|
||||||
"git.g3e.fr/syonad/two/pkg/worker"
|
"git.g3e.fr/syonad/two/pkg/worker"
|
||||||
|
|
@ -17,20 +19,32 @@ type Dispatcher struct {
|
||||||
queue *worker.Queue
|
queue *worker.Queue
|
||||||
db *badger.DB
|
db *badger.DB
|
||||||
cfg *configuration.Config
|
cfg *configuration.Config
|
||||||
|
logger *slog.Logger
|
||||||
}
|
}
|
||||||
|
|
||||||
func New(queue *worker.Queue, db *badger.DB, cfg *configuration.Config) *Dispatcher {
|
func New(queue *worker.Queue, db *badger.DB, cfg *configuration.Config, logger *slog.Logger) *Dispatcher {
|
||||||
return &Dispatcher{queue: queue, db: db, cfg: cfg}
|
return &Dispatcher{queue: queue, db: db, cfg: cfg, logger: logger}
|
||||||
}
|
}
|
||||||
|
|
||||||
func (d *Dispatcher) Prepare(cmd Command) error {
|
func (d *Dispatcher) Prepare(cmd Command) error {
|
||||||
|
d.logger.Debug("prepare", "command", fmt.Sprintf("%T", cmd))
|
||||||
return cmd.Prepare(d.db, d.cfg)
|
return cmd.Prepare(d.db, d.cfg)
|
||||||
}
|
}
|
||||||
|
|
||||||
func (d *Dispatcher) Dispatch(cmd Command) {
|
func (d *Dispatcher) Dispatch(cmd Command) {
|
||||||
|
cmdType := fmt.Sprintf("%T", cmd)
|
||||||
|
d.logger.Debug("dispatch", "command", cmdType)
|
||||||
d.queue.Submit(func() {
|
d.queue.Submit(func() {
|
||||||
if err := cmd.Execute(d.db, d.cfg); err != nil {
|
start := time.Now()
|
||||||
log.Printf("command error (%T): %v", cmd, err)
|
err := cmd.Execute(d.db, d.cfg)
|
||||||
|
attrs := []any{
|
||||||
|
"command", cmdType,
|
||||||
|
"duration_ms", time.Since(start).Milliseconds(),
|
||||||
|
}
|
||||||
|
if err != nil {
|
||||||
|
d.logger.Error("command failed", append(attrs, "error", err)...)
|
||||||
|
} else {
|
||||||
|
d.logger.Info("command done", attrs...)
|
||||||
}
|
}
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -2,7 +2,6 @@ package dispatcher
|
||||||
|
|
||||||
import (
|
import (
|
||||||
"fmt"
|
"fmt"
|
||||||
"os"
|
|
||||||
"strconv"
|
"strconv"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
|
|
@ -77,13 +76,13 @@ func (c DeleteSubnetCommand) Prepare(db *badger.DB, _ *configuration.Config) err
|
||||||
|
|
||||||
func (c DeleteSubnetCommand) Execute(db *badger.DB, _ *configuration.Config) error {
|
func (c DeleteSubnetCommand) Execute(db *badger.DB, _ *configuration.Config) error {
|
||||||
if err := subnet.DeleteSubnet(db, c.Name); err != nil {
|
if err := subnet.DeleteSubnet(db, c.Name); err != nil {
|
||||||
fmt.Println(err)
|
return err
|
||||||
os.Exit(1)
|
|
||||||
}
|
}
|
||||||
if state, err := kv.GetFromDB(db, "subnet/"+c.Name+"/state"); err != nil {
|
state, err := kv.GetFromDB(db, "subnet/"+c.Name+"/state")
|
||||||
fmt.Println(err)
|
if err != nil {
|
||||||
os.Exit(1)
|
return err
|
||||||
} else if state == "deleted" {
|
}
|
||||||
|
if state == "deleted" {
|
||||||
kv.DeleteInDB(db, "subnet/"+c.Name)
|
kv.DeleteInDB(db, "subnet/"+c.Name)
|
||||||
}
|
}
|
||||||
return nil
|
return nil
|
||||||
|
|
|
||||||
25
pkg/logger/logger.go
Normal file
25
pkg/logger/logger.go
Normal file
|
|
@ -0,0 +1,25 @@
|
||||||
|
package logger
|
||||||
|
|
||||||
|
import (
|
||||||
|
"log/slog"
|
||||||
|
"os"
|
||||||
|
)
|
||||||
|
|
||||||
|
var Level = new(slog.LevelVar)
|
||||||
|
|
||||||
|
func New(level string, debug bool) *slog.Logger {
|
||||||
|
switch level {
|
||||||
|
case "debug":
|
||||||
|
Level.Set(slog.LevelDebug)
|
||||||
|
case "warn":
|
||||||
|
Level.Set(slog.LevelWarn)
|
||||||
|
case "error":
|
||||||
|
Level.Set(slog.LevelError)
|
||||||
|
default:
|
||||||
|
Level.Set(slog.LevelInfo)
|
||||||
|
}
|
||||||
|
if debug {
|
||||||
|
Level.Set(slog.LevelDebug)
|
||||||
|
}
|
||||||
|
return slog.New(slog.NewJSONHandler(os.Stderr, &slog.HandlerOptions{Level: Level}))
|
||||||
|
}
|
||||||
Loading…
Add table
Add a link
Reference in a new issue