f-29: stop: add a proper stop
All checks were successful
Release Pipeline / upload-assets (metadata@.service, systemd/metadata@.service) (push) Successful in 4s
Release Pipeline / upload-assets (run-dnsmasq-in-netns.sh, scripts/run-dnsmasq-in-netns.sh) (push) Successful in 4s
Release Pipeline / build (agent, amd64, linux) (push) Successful in 0s
Release Pipeline / build (metadata, amd64, linux) (push) Successful in 0s
Release Pipeline / checksums (push) Successful in 4s
Release Pipeline / set-release-target (push) Successful in 2s
Release Pipeline / upload-assets (agent.service, systemd/agent.service) (push) Successful in 4s
Release Pipeline / upload-assets (dnsmasq@.service, systemd/dnsmasq@.service) (push) Successful in 4s
Release Pipeline / release (push) Successful in 11s
Release Pipeline / publish (push) Successful in 0s
Release Pipeline / build (push) Successful in 1m30s
All checks were successful
Release Pipeline / upload-assets (metadata@.service, systemd/metadata@.service) (push) Successful in 4s
Release Pipeline / upload-assets (run-dnsmasq-in-netns.sh, scripts/run-dnsmasq-in-netns.sh) (push) Successful in 4s
Release Pipeline / build (agent, amd64, linux) (push) Successful in 0s
Release Pipeline / build (metadata, amd64, linux) (push) Successful in 0s
Release Pipeline / checksums (push) Successful in 4s
Release Pipeline / set-release-target (push) Successful in 2s
Release Pipeline / upload-assets (agent.service, systemd/agent.service) (push) Successful in 4s
Release Pipeline / upload-assets (dnsmasq@.service, systemd/dnsmasq@.service) (push) Successful in 4s
Release Pipeline / release (push) Successful in 11s
Release Pipeline / publish (push) Successful in 0s
Release Pipeline / build (push) Successful in 1m30s
Signed-off-by: GnomeZworc <nicolas.boufidjeline@g3e.fr>
This commit is contained in:
parent
1ce84ffd58
commit
cdbb5012b1
9 changed files with 333 additions and 29 deletions
|
|
@ -5,6 +5,9 @@ import (
|
||||||
"flag"
|
"flag"
|
||||||
"fmt"
|
"fmt"
|
||||||
"log/slog"
|
"log/slog"
|
||||||
|
"os"
|
||||||
|
"os/signal"
|
||||||
|
"syscall"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
agentapi "git.g3e.fr/syonad/two/internal/api/agent"
|
agentapi "git.g3e.fr/syonad/two/internal/api/agent"
|
||||||
|
|
@ -21,6 +24,8 @@ import (
|
||||||
"github.com/prometheus/client_golang/prometheus"
|
"github.com/prometheus/client_golang/prometheus"
|
||||||
)
|
)
|
||||||
|
|
||||||
|
const shutdownTimeout = 20 * time.Second
|
||||||
|
|
||||||
func main() {
|
func main() {
|
||||||
confFile := flag.String("config", "/etc/two/agent.yml", "config file path")
|
confFile := flag.String("config", "/etc/two/agent.yml", "config file path")
|
||||||
flag.Parse()
|
flag.Parse()
|
||||||
|
|
@ -34,7 +39,12 @@ func main() {
|
||||||
log := logger.New(cfg.Logger.Level, cfg.Logger.Debug)
|
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()
|
closeDB := true
|
||||||
|
defer func() {
|
||||||
|
if closeDB {
|
||||||
|
db.Close()
|
||||||
|
}
|
||||||
|
}()
|
||||||
|
|
||||||
// Avant tout démarrage de service : la DB peut porter l'ancien vocabulaire
|
// Avant tout démarrage de service : la DB peut porter l'ancien vocabulaire
|
||||||
// d'états, et des ressources transitoires orphelines d'un arrêt précédent.
|
// d'états, et des ressources transitoires orphelines d'un arrêt précédent.
|
||||||
|
|
@ -60,19 +70,69 @@ func main() {
|
||||||
"debug", cfg.Logger.Debug,
|
"debug", cfg.Logger.Debug,
|
||||||
)
|
)
|
||||||
|
|
||||||
|
ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM)
|
||||||
|
defer stop()
|
||||||
|
|
||||||
d := dispatcher.New(q, db, cfg, log.With(slog.String("component", "dispatcher")))
|
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)
|
apiSrv := agentapi.New(d, db, log.With(slog.String("component", "api")), apiAddr)
|
||||||
|
go apiSrv.Start()
|
||||||
|
|
||||||
|
promSrv := promserver.New(promAddr, registry)
|
||||||
|
go promSrv.Start()
|
||||||
|
|
||||||
|
var adminSrv *kv.AdminServer
|
||||||
if cfg.Admin.Enabled {
|
if cfg.Admin.Enabled {
|
||||||
adminAddr := fmt.Sprintf("%s:%d", cfg.Admin.Address, cfg.Admin.Port)
|
adminAddr := fmt.Sprintf("%s:%d", cfg.Admin.Address, cfg.Admin.Port)
|
||||||
go kv.NewAdminServer(db, log.With(slog.String("component", "admin"))).Start(adminAddr)
|
adminSrv = kv.NewAdminServer(db, log.With(slog.String("component", "admin")), adminAddr)
|
||||||
|
go adminSrv.Start()
|
||||||
}
|
}
|
||||||
|
|
||||||
if cfg.Watchdog.Enabled {
|
if cfg.Watchdog.Enabled {
|
||||||
wlog := log.With(slog.String("component", "watchdog"))
|
wlog := log.With(slog.String("component", "watchdog"))
|
||||||
go watchdog.New(db, cfg, notify.NewStderr(wlog), wlog,
|
go watchdog.New(db, cfg, notify.NewStderr(wlog), wlog,
|
||||||
time.Duration(cfg.Watchdog.IntervalSeconds)*time.Second,
|
time.Duration(cfg.Watchdog.IntervalSeconds)*time.Second,
|
||||||
).Run(context.Background())
|
).Run(ctx)
|
||||||
}
|
}
|
||||||
|
|
||||||
select {}
|
<-ctx.Done()
|
||||||
|
stop()
|
||||||
|
|
||||||
|
servers := map[string]httpShutdowner{"api": apiSrv, "prometheus": promSrv}
|
||||||
|
if adminSrv != nil {
|
||||||
|
servers["admin"] = adminSrv
|
||||||
|
}
|
||||||
|
closeDB = shutdown(log, q, servers, shutdownTimeout)
|
||||||
|
}
|
||||||
|
|
||||||
|
type httpShutdowner interface {
|
||||||
|
Shutdown(context.Context) error
|
||||||
|
}
|
||||||
|
|
||||||
|
func shutdown(log *slog.Logger, q *worker.Queue, servers map[string]httpShutdowner, timeout time.Duration) bool {
|
||||||
|
log.Info("shutting down", "timeout", timeout)
|
||||||
|
|
||||||
|
ctx, cancel := context.WithTimeout(context.Background(), timeout)
|
||||||
|
defer cancel()
|
||||||
|
|
||||||
|
for name, srv := range servers {
|
||||||
|
if err := srv.Shutdown(ctx); err != nil {
|
||||||
|
log.Error("http server shutdown", "server", name, "error", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
drained := make(chan struct{})
|
||||||
|
go func() {
|
||||||
|
q.Stop()
|
||||||
|
close(drained)
|
||||||
|
}()
|
||||||
|
|
||||||
|
select {
|
||||||
|
case <-drained:
|
||||||
|
log.Info("workers drained")
|
||||||
|
return true
|
||||||
|
case <-ctx.Done():
|
||||||
|
log.Error("workers still running after timeout, leaving database untouched")
|
||||||
|
return false
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
98
cmd/agent/main_test.go
Normal file
98
cmd/agent/main_test.go
Normal file
|
|
@ -0,0 +1,98 @@
|
||||||
|
package main
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"errors"
|
||||||
|
"io"
|
||||||
|
"log/slog"
|
||||||
|
"strings"
|
||||||
|
"sync/atomic"
|
||||||
|
"testing"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"git.g3e.fr/syonad/two/pkg/worker"
|
||||||
|
)
|
||||||
|
|
||||||
|
type fakeServer struct {
|
||||||
|
called atomic.Bool
|
||||||
|
err error
|
||||||
|
}
|
||||||
|
|
||||||
|
func (f *fakeServer) Shutdown(context.Context) error {
|
||||||
|
f.called.Store(true)
|
||||||
|
return f.err
|
||||||
|
}
|
||||||
|
|
||||||
|
func discardLogger() *slog.Logger {
|
||||||
|
return slog.New(slog.NewTextHandler(io.Discard, nil))
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestShutdown_DrainReussi(t *testing.T) {
|
||||||
|
q := worker.New(10)
|
||||||
|
q.Start(2)
|
||||||
|
|
||||||
|
var done atomic.Int32
|
||||||
|
for range 3 {
|
||||||
|
q.Submit(func() {
|
||||||
|
time.Sleep(20 * time.Millisecond)
|
||||||
|
done.Add(1)
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
api, prom := &fakeServer{}, &fakeServer{}
|
||||||
|
servers := map[string]httpShutdowner{"api": api, "prometheus": prom}
|
||||||
|
|
||||||
|
if !shutdown(discardLogger(), q, servers, 5*time.Second) {
|
||||||
|
t.Fatal("un drainage réussi doit autoriser la fermeture de la base")
|
||||||
|
}
|
||||||
|
if !api.called.Load() || !prom.called.Load() {
|
||||||
|
t.Error("tous les serveurs HTTP doivent être arrêtés")
|
||||||
|
}
|
||||||
|
if got := done.Load(); got != 3 {
|
||||||
|
t.Errorf("les 3 tâches devaient se terminer, %d terminées", got)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestShutdown_TimeoutLaisseLaBaseIntacte(t *testing.T) {
|
||||||
|
q := worker.New(10)
|
||||||
|
q.Start(1)
|
||||||
|
q.Submit(func() { time.Sleep(2 * time.Second) })
|
||||||
|
|
||||||
|
var buf strings.Builder
|
||||||
|
log := slog.New(slog.NewTextHandler(&buf, nil))
|
||||||
|
|
||||||
|
if shutdown(log, q, map[string]httpShutdowner{}, 50*time.Millisecond) {
|
||||||
|
t.Fatal("un drainage incomplet ne doit pas autoriser la fermeture de la base")
|
||||||
|
}
|
||||||
|
if !strings.Contains(buf.String(), "leaving database untouched") {
|
||||||
|
t.Errorf("le dépassement devrait être logué, obtenu %q", buf.String())
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestShutdown_ErreurServeurNEmpechePasLeDrainage(t *testing.T) {
|
||||||
|
q := worker.New(10)
|
||||||
|
q.Start(1)
|
||||||
|
|
||||||
|
var buf strings.Builder
|
||||||
|
log := slog.New(slog.NewTextHandler(&buf, nil))
|
||||||
|
servers := map[string]httpShutdowner{
|
||||||
|
"api": &fakeServer{err: errors.New("boom")},
|
||||||
|
"prometheus": &fakeServer{},
|
||||||
|
}
|
||||||
|
|
||||||
|
if !shutdown(log, q, servers, 5*time.Second) {
|
||||||
|
t.Fatal("une erreur d'arrêt HTTP ne doit pas empêcher le drainage")
|
||||||
|
}
|
||||||
|
if !strings.Contains(buf.String(), "http server shutdown") {
|
||||||
|
t.Errorf("l'erreur devrait être loguée, obtenu %q", buf.String())
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestShutdown_SansServeur(t *testing.T) {
|
||||||
|
q := worker.New(10)
|
||||||
|
q.Start(1)
|
||||||
|
|
||||||
|
if !shutdown(discardLogger(), q, map[string]httpShutdowner{}, 5*time.Second) {
|
||||||
|
t.Fatal("l'absence de serveur ne doit pas empêcher un arrêt propre")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
@ -23,5 +23,5 @@ func newTestServer(t *testing.T) (*Server, *badger.DB) {
|
||||||
cfg := &configuration.Config{DefaultInterface: "br-test"}
|
cfg := &configuration.Config{DefaultInterface: "br-test"}
|
||||||
logger := slog.New(slog.NewTextHandler(io.Discard, nil))
|
logger := slog.New(slog.NewTextHandler(io.Discard, nil))
|
||||||
d := dispatcher.New(q, db, cfg, logger)
|
d := dispatcher.New(q, db, cfg, logger)
|
||||||
return New(d, db, logger), db
|
return New(d, db, logger, "127.0.0.1:0"), db
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -1,8 +1,10 @@
|
||||||
package agentapi
|
package agentapi
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"context"
|
||||||
"crypto/rand"
|
"crypto/rand"
|
||||||
"encoding/hex"
|
"encoding/hex"
|
||||||
|
"errors"
|
||||||
"log/slog"
|
"log/slog"
|
||||||
"net/http"
|
"net/http"
|
||||||
"time"
|
"time"
|
||||||
|
|
@ -15,13 +17,11 @@ type Server struct {
|
||||||
dispatcher *dispatcher.Dispatcher
|
dispatcher *dispatcher.Dispatcher
|
||||||
db *badger.DB
|
db *badger.DB
|
||||||
logger *slog.Logger
|
logger *slog.Logger
|
||||||
|
srv *http.Server
|
||||||
}
|
}
|
||||||
|
|
||||||
func New(d *dispatcher.Dispatcher, db *badger.DB, logger *slog.Logger) *Server {
|
func New(d *dispatcher.Dispatcher, db *badger.DB, logger *slog.Logger, address string) *Server {
|
||||||
return &Server{dispatcher: d, db: db, logger: logger}
|
s := &Server{dispatcher: d, db: db, logger: logger}
|
||||||
}
|
|
||||||
|
|
||||||
func (s *Server) Start(address string) {
|
|
||||||
mux := http.NewServeMux()
|
mux := http.NewServeMux()
|
||||||
mux.HandleFunc("/vpcs", s.VpcsHandler)
|
mux.HandleFunc("/vpcs", s.VpcsHandler)
|
||||||
mux.HandleFunc("/vpcs/", s.VpcByNameHandler)
|
mux.HandleFunc("/vpcs/", s.VpcByNameHandler)
|
||||||
|
|
@ -29,12 +29,21 @@ func (s *Server) Start(address string) {
|
||||||
mux.HandleFunc("/subnets/", s.SubnetByNameHandler)
|
mux.HandleFunc("/subnets/", s.SubnetByNameHandler)
|
||||||
mux.HandleFunc("/vms", s.VmsHandler)
|
mux.HandleFunc("/vms", s.VmsHandler)
|
||||||
mux.HandleFunc("/vms/", s.VmByNameHandler)
|
mux.HandleFunc("/vms/", s.VmByNameHandler)
|
||||||
s.logger.Info("API server listening", "address", address)
|
s.srv = &http.Server{Addr: address, Handler: s.logMiddleware(mux)}
|
||||||
if err := http.ListenAndServe(address, s.logMiddleware(mux)); err != nil {
|
return s
|
||||||
|
}
|
||||||
|
|
||||||
|
func (s *Server) Start() {
|
||||||
|
s.logger.Info("API server listening", "address", s.srv.Addr)
|
||||||
|
if err := s.srv.ListenAndServe(); err != nil && !errors.Is(err, http.ErrServerClosed) {
|
||||||
s.logger.Error("API server stopped", "error", err)
|
s.logger.Error("API server stopped", "error", err)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (s *Server) Shutdown(ctx context.Context) error {
|
||||||
|
return s.srv.Shutdown(ctx)
|
||||||
|
}
|
||||||
|
|
||||||
type statusWriter struct {
|
type statusWriter struct {
|
||||||
http.ResponseWriter
|
http.ResponseWriter
|
||||||
status int
|
status int
|
||||||
|
|
|
||||||
|
|
@ -1,6 +1,8 @@
|
||||||
package kv
|
package kv
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"context"
|
||||||
|
"errors"
|
||||||
"fmt"
|
"fmt"
|
||||||
"log/slog"
|
"log/slog"
|
||||||
"net/http"
|
"net/http"
|
||||||
|
|
@ -12,21 +14,28 @@ import (
|
||||||
type AdminServer struct {
|
type AdminServer struct {
|
||||||
db *badger.DB
|
db *badger.DB
|
||||||
logger *slog.Logger
|
logger *slog.Logger
|
||||||
|
srv *http.Server
|
||||||
}
|
}
|
||||||
|
|
||||||
func NewAdminServer(db *badger.DB, logger *slog.Logger) *AdminServer {
|
func NewAdminServer(db *badger.DB, logger *slog.Logger, address string) *AdminServer {
|
||||||
return &AdminServer{db: db, logger: logger}
|
s := &AdminServer{db: db, logger: logger}
|
||||||
}
|
|
||||||
|
|
||||||
func (s *AdminServer) Start(address string) {
|
|
||||||
mux := http.NewServeMux()
|
mux := http.NewServeMux()
|
||||||
mux.HandleFunc("/db", s.dbHandler)
|
mux.HandleFunc("/db", s.dbHandler)
|
||||||
s.logger.Info("admin server listening", "address", address)
|
s.srv = &http.Server{Addr: address, Handler: mux}
|
||||||
if err := http.ListenAndServe(address, mux); err != nil {
|
return s
|
||||||
|
}
|
||||||
|
|
||||||
|
func (s *AdminServer) Start() {
|
||||||
|
s.logger.Info("admin server listening", "address", s.srv.Addr)
|
||||||
|
if err := s.srv.ListenAndServe(); err != nil && !errors.Is(err, http.ErrServerClosed) {
|
||||||
s.logger.Error("admin server stopped", "error", err)
|
s.logger.Error("admin server stopped", "error", err)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (s *AdminServer) Shutdown(ctx context.Context) error {
|
||||||
|
return s.srv.Shutdown(ctx)
|
||||||
|
}
|
||||||
|
|
||||||
func (s *AdminServer) dbHandler(w http.ResponseWriter, r *http.Request) {
|
func (s *AdminServer) dbHandler(w http.ResponseWriter, r *http.Request) {
|
||||||
if r.Method != http.MethodGet {
|
if r.Method != http.MethodGet {
|
||||||
http.Error(w, "method not allowed", http.StatusMethodNotAllowed)
|
http.Error(w, "method not allowed", http.StatusMethodNotAllowed)
|
||||||
|
|
|
||||||
|
|
@ -1,6 +1,8 @@
|
||||||
package promserver
|
package promserver
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"context"
|
||||||
|
"errors"
|
||||||
"log"
|
"log"
|
||||||
"net/http"
|
"net/http"
|
||||||
|
|
||||||
|
|
@ -8,13 +10,30 @@ import (
|
||||||
"github.com/prometheus/client_golang/prometheus/promhttp"
|
"github.com/prometheus/client_golang/prometheus/promhttp"
|
||||||
)
|
)
|
||||||
|
|
||||||
// Start launches the Prometheus metrics HTTP server on the given address.
|
// Server exposes the Prometheus metrics endpoint.
|
||||||
// The provided registry is used to expose metrics at /metrics.
|
type Server struct {
|
||||||
func Start(address string, registry *prometheus.Registry) {
|
srv *http.Server
|
||||||
|
}
|
||||||
|
|
||||||
|
// New builds the metrics server for the given address and registry.
|
||||||
|
func New(address string, registry *prometheus.Registry) *Server {
|
||||||
mux := http.NewServeMux()
|
mux := http.NewServeMux()
|
||||||
mux.Handle("/metrics", promhttp.HandlerFor(registry, promhttp.HandlerOpts{
|
mux.Handle("/metrics", promhttp.HandlerFor(registry, promhttp.HandlerOpts{
|
||||||
EnableOpenMetrics: true,
|
EnableOpenMetrics: true,
|
||||||
}))
|
}))
|
||||||
log.Printf("Prometheus server listening on %s", address)
|
return &Server{srv: &http.Server{Addr: address, Handler: mux}}
|
||||||
log.Fatal(http.ListenAndServe(address, mux))
|
}
|
||||||
|
|
||||||
|
// Start blocks until the server stops. A failure is logged, never fatal: an
|
||||||
|
// unavailable metrics endpoint must not bring the whole agent down.
|
||||||
|
func (s *Server) Start() {
|
||||||
|
log.Printf("Prometheus server listening on %s", s.srv.Addr)
|
||||||
|
if err := s.srv.ListenAndServe(); err != nil && !errors.Is(err, http.ErrServerClosed) {
|
||||||
|
log.Printf("Prometheus server stopped: %v", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// Shutdown stops the server, waiting for in-flight requests until ctx expires.
|
||||||
|
func (s *Server) Shutdown(ctx context.Context) error {
|
||||||
|
return s.srv.Shutdown(ctx)
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -1,13 +1,19 @@
|
||||||
package worker
|
package worker
|
||||||
|
|
||||||
import "log"
|
import (
|
||||||
|
"log"
|
||||||
|
"sync"
|
||||||
|
)
|
||||||
|
|
||||||
// Task is a function to be executed asynchronously by a worker.
|
// Task is a function to be executed asynchronously by a worker.
|
||||||
type Task func()
|
type Task func()
|
||||||
|
|
||||||
// Queue is a FIFO channel-backed task queue consumed by worker goroutines.
|
// Queue is a FIFO channel-backed task queue consumed by worker goroutines.
|
||||||
type Queue struct {
|
type Queue struct {
|
||||||
tasks chan Task
|
tasks chan Task
|
||||||
|
wg sync.WaitGroup
|
||||||
|
mu sync.RWMutex
|
||||||
|
stopped bool
|
||||||
}
|
}
|
||||||
|
|
||||||
// New creates a Queue with the given channel buffer size.
|
// New creates a Queue with the given channel buffer size.
|
||||||
|
|
@ -15,8 +21,16 @@ func New(bufferSize int) *Queue {
|
||||||
return &Queue{tasks: make(chan Task, bufferSize)}
|
return &Queue{tasks: make(chan Task, bufferSize)}
|
||||||
}
|
}
|
||||||
|
|
||||||
// Submit enqueues a task. Blocks if the queue is full.
|
// Submit enqueues a task. Blocks if the queue is full. Tasks submitted after
|
||||||
|
// Stop are rejected and logged rather than enqueued.
|
||||||
func (q *Queue) Submit(t Task) {
|
func (q *Queue) Submit(t Task) {
|
||||||
|
q.mu.RLock()
|
||||||
|
defer q.mu.RUnlock()
|
||||||
|
|
||||||
|
if q.stopped {
|
||||||
|
log.Print("worker: queue stopped, task rejected")
|
||||||
|
return
|
||||||
|
}
|
||||||
q.tasks <- t
|
q.tasks <- t
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -24,10 +38,28 @@ func (q *Queue) Submit(t Task) {
|
||||||
func (q *Queue) Start(n int) {
|
func (q *Queue) Start(n int) {
|
||||||
log.Printf("worker: starting %d workers", n)
|
log.Printf("worker: starting %d workers", n)
|
||||||
for i := range n {
|
for i := range n {
|
||||||
|
q.wg.Add(1)
|
||||||
go func(id int) {
|
go func(id int) {
|
||||||
|
defer q.wg.Done()
|
||||||
for task := range q.tasks {
|
for task := range q.tasks {
|
||||||
task()
|
task()
|
||||||
}
|
}
|
||||||
}(i)
|
}(i)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Stop rejects new tasks, then waits for the queued and in-flight ones to
|
||||||
|
// finish. The write lock is what makes closing the channel safe: it is only
|
||||||
|
// taken once every in-flight Submit has released its read lock.
|
||||||
|
func (q *Queue) Stop() {
|
||||||
|
q.mu.Lock()
|
||||||
|
if q.stopped {
|
||||||
|
q.mu.Unlock()
|
||||||
|
return
|
||||||
|
}
|
||||||
|
q.stopped = true
|
||||||
|
close(q.tasks)
|
||||||
|
q.mu.Unlock()
|
||||||
|
|
||||||
|
q.wg.Wait()
|
||||||
|
}
|
||||||
|
|
|
||||||
|
|
@ -101,3 +101,79 @@ func TestQueue_SubmitBlocksWhenFull(t *testing.T) {
|
||||||
t.Error("Submit aurait dû se débloquer après démarrage d'un worker")
|
t.Error("Submit aurait dû se débloquer après démarrage d'un worker")
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func TestStop_AttendLesTachesEnCours(t *testing.T) {
|
||||||
|
q := New(10)
|
||||||
|
q.Start(2)
|
||||||
|
|
||||||
|
var mu sync.Mutex
|
||||||
|
done := 0
|
||||||
|
for range 5 {
|
||||||
|
q.Submit(func() {
|
||||||
|
time.Sleep(20 * time.Millisecond)
|
||||||
|
mu.Lock()
|
||||||
|
done++
|
||||||
|
mu.Unlock()
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
q.Stop()
|
||||||
|
|
||||||
|
mu.Lock()
|
||||||
|
defer mu.Unlock()
|
||||||
|
if done != 5 {
|
||||||
|
t.Errorf("Stop devrait attendre les 5 tâches, %d terminées", done)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestStop_RejetteLesTachesSuivantes(t *testing.T) {
|
||||||
|
q := New(10)
|
||||||
|
q.Start(1)
|
||||||
|
q.Stop()
|
||||||
|
|
||||||
|
executed := make(chan struct{}, 1)
|
||||||
|
q.Submit(func() { executed <- struct{}{} })
|
||||||
|
|
||||||
|
select {
|
||||||
|
case <-executed:
|
||||||
|
t.Error("une tâche soumise après Stop ne doit pas être exécutée")
|
||||||
|
case <-time.After(100 * time.Millisecond):
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestStop_Idempotent(t *testing.T) {
|
||||||
|
q := New(10)
|
||||||
|
q.Start(1)
|
||||||
|
q.Stop()
|
||||||
|
q.Stop()
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestStop_SansTacheEnCours(t *testing.T) {
|
||||||
|
q := New(10)
|
||||||
|
q.Start(3)
|
||||||
|
|
||||||
|
done := make(chan struct{})
|
||||||
|
go func() { q.Stop(); close(done) }()
|
||||||
|
|
||||||
|
select {
|
||||||
|
case <-done:
|
||||||
|
case <-time.After(2 * time.Second):
|
||||||
|
t.Fatal("Stop n'a pas rendu la main sur une file vide")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestSubmit_ConcurrentAvecStop(t *testing.T) {
|
||||||
|
q := New(100)
|
||||||
|
q.Start(4)
|
||||||
|
|
||||||
|
var wg sync.WaitGroup
|
||||||
|
for range 20 {
|
||||||
|
wg.Add(1)
|
||||||
|
go func() {
|
||||||
|
defer wg.Done()
|
||||||
|
q.Submit(func() {})
|
||||||
|
}()
|
||||||
|
}
|
||||||
|
q.Stop()
|
||||||
|
wg.Wait()
|
||||||
|
}
|
||||||
|
|
|
||||||
|
|
@ -5,6 +5,7 @@ After=network.target
|
||||||
[Service]
|
[Service]
|
||||||
Type=simple
|
Type=simple
|
||||||
ExecStart=/opt/two/bin/agent
|
ExecStart=/opt/two/bin/agent
|
||||||
|
TimeoutStopSec=30
|
||||||
|
|
||||||
[Install]
|
[Install]
|
||||||
WantedBy=multi-user.target
|
WantedBy=multi-user.target
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue