Compare commits

..

No commits in common. "921a5ca96e8f48a24c76e8174b2bbc45f306609c" and "648a64782c4e941edd488ccb2b4be4b710ff1545" have entirely different histories.

7 changed files with 29 additions and 122 deletions

View file

@ -3,14 +3,13 @@ package main
import ( import (
"flag" "flag"
"fmt" "fmt"
"log/slog" "log"
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"
@ -22,12 +21,9 @@ func main() {
cfg, err := configuration.LoadConfig(*confFile) cfg, err := configuration.LoadConfig(*confFile)
if err != nil { if err != nil {
slog.Error("failed to load config", "error", err) log.Fatalf("failed to load config: %v", 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()
@ -40,16 +36,8 @@ 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)
log.Info("starting agent", d := dispatcher.New(q, db, cfg)
"api", apiAddr, go agentapi.New(d, db).Start(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 {}

View file

@ -34,10 +34,3 @@ 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

View file

@ -1,11 +1,8 @@
package agentapi package agentapi
import ( import (
"crypto/rand" "log"
"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"
@ -14,11 +11,10 @@ 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, logger *slog.Logger) *Server { func New(d *dispatcher.Dispatcher, db *badger.DB) *Server {
return &Server{dispatcher: d, db: db, logger: logger} return &Server{dispatcher: d, db: db}
} }
func (s *Server) Start(address string) { func (s *Server) Start(address string) {
@ -27,39 +23,13 @@ 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)
s.logger.Info("API server listening", "address", address) log.Printf("API server listening on %s", address)
if err := http.ListenAndServe(address, s.logMiddleware(mux)); err != nil { log.Fatal(http.ListenAndServe(address, logMiddleware(mux)))
s.logger.Error("API server stopped", "error", err)
}
} }
type statusWriter struct { func logMiddleware(next http.Handler) http.Handler {
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) {
var b [4]byte log.Printf("%s %s %s", r.RemoteAddr, r.Method, r.URL.Path)
rand.Read(b[:]) next.ServeHTTP(w, r)
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,
)
}) })
} }

View file

@ -21,13 +21,9 @@ type Config struct {
BufferSize int `mapstructure:"buffer_size"` BufferSize int `mapstructure:"buffer_size"`
} `mapstructure:"worker"` } `mapstructure:"worker"`
Dispatcher struct { Dispatcher 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"`
} }
@ -47,8 +43,6 @@ 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()

View file

@ -1,9 +1,7 @@
package dispatcher package dispatcher
import ( import (
"fmt" "log"
"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"
@ -16,35 +14,23 @@ type Command interface {
} }
type Dispatcher struct { 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, logger *slog.Logger) *Dispatcher { func New(queue *worker.Queue, db *badger.DB, cfg *configuration.Config) *Dispatcher {
return &Dispatcher{queue: queue, db: db, cfg: cfg, logger: logger} return &Dispatcher{queue: queue, db: db, cfg: cfg}
} }
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() {
start := time.Now() if err := cmd.Execute(d.db, d.cfg); err != nil {
err := cmd.Execute(d.db, d.cfg) log.Printf("command error (%T): %v", cmd, err)
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...)
} }
}) })
} }

View file

@ -2,6 +2,7 @@ package dispatcher
import ( import (
"fmt" "fmt"
"os"
"strconv" "strconv"
"time" "time"
@ -76,13 +77,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 {
return err fmt.Println(err)
os.Exit(1)
} }
state, err := kv.GetFromDB(db, "subnet/"+c.Name+"/state") if state, err := kv.GetFromDB(db, "subnet/"+c.Name+"/state"); err != nil {
if err != nil { fmt.Println(err)
return err os.Exit(1)
} } else if state == "deleted" {
if state == "deleted" {
kv.DeleteInDB(db, "subnet/"+c.Name) kv.DeleteInDB(db, "subnet/"+c.Name)
} }
return nil return nil

View file

@ -1,25 +0,0 @@
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}))
}