From 2ce7ebafe727d7efd3fce00ee573fe10acbb830b Mon Sep 17 00:00:00 2001 From: GnomeZworc Date: Mon, 17 Aug 2026 23:52:48 +0200 Subject: [PATCH 1/4] f-29: watchdog: add specific config #29 Signed-off-by: GnomeZworc --- conf/agent/config.exemple.yml | 7 +++ internal/config/agent/config_test.go | 65 ++++++++++++++++++++++++++++ internal/config/agent/struct.go | 6 +++ 3 files changed, 78 insertions(+) diff --git a/conf/agent/config.exemple.yml b/conf/agent/config.exemple.yml index 0517eac..72e88f8 100644 --- a/conf/agent/config.exemple.yml +++ b/conf/agent/config.exemple.yml @@ -51,6 +51,13 @@ qemu: monitor_dir: "/run/two/vms/monitor" qmp_dir: "/run/two/vms/qmp" +# Consistency watchdog: periodically checks that resources marked running in the +# database still exist on the system, and reports the gaps. Read-only, never repairs. +# A persisting gap is reported at every tick (no deduplication). +watchdog: + enabled: true + interval_seconds: 60 + # Admin API (read-only DB inspection, loopback only) admin: enabled: false diff --git a/internal/config/agent/config_test.go b/internal/config/agent/config_test.go index d0f3d31..a4d5b22 100644 --- a/internal/config/agent/config_test.go +++ b/internal/config/agent/config_test.go @@ -81,3 +81,68 @@ database: t.Errorf("attendu %q, obtenu %q", "/opt/two/data", cfg.Database.Path) } } + +func TestLoadConfig_WatchdogDefauts(t *testing.T) { + path := writeYAML(t, "") + cfg, err := LoadConfig(path) + if err != nil { + t.Fatalf("LoadConfig a échoué : %v", err) + } + if cfg.Watchdog.Enabled { + t.Error("watchdog.enabled devrait être false par défaut") + } + if cfg.Watchdog.IntervalSeconds != 60 { + t.Errorf("watchdog.interval_seconds attendu 60, obtenu %d", cfg.Watchdog.IntervalSeconds) + } +} + +func TestLoadConfig_WatchdogValeursExplicites(t *testing.T) { + path := writeYAML(t, ` +watchdog: + enabled: true + interval_seconds: 30 +`) + cfg, err := LoadConfig(path) + if err != nil { + t.Fatalf("LoadConfig a échoué : %v", err) + } + if !cfg.Watchdog.Enabled { + t.Error("watchdog.enabled attendu true") + } + if cfg.Watchdog.IntervalSeconds != 30 { + t.Errorf("watchdog.interval_seconds attendu 30, obtenu %d", cfg.Watchdog.IntervalSeconds) + } +} + +func TestLoadConfig_WatchdogActiveSansIntervalle(t *testing.T) { + path := writeYAML(t, ` +watchdog: + enabled: true +`) + cfg, err := LoadConfig(path) + if err != nil { + t.Fatalf("LoadConfig a échoué : %v", err) + } + if !cfg.Watchdog.Enabled { + t.Error("watchdog.enabled attendu true") + } + if cfg.Watchdog.IntervalSeconds != 60 { + t.Errorf("watchdog.interval_seconds attendu 60 (défaut viper), obtenu %d", cfg.Watchdog.IntervalSeconds) + } +} + +func TestLoadConfig_ExempleFourniEstValide(t *testing.T) { + cfg, err := LoadConfig("../../../conf/agent/config.exemple.yml") + if err != nil { + t.Fatalf("config.exemple.yml illisible : %v", err) + } + if !cfg.Watchdog.Enabled { + t.Error("config.exemple.yml devrait activer le watchdog") + } + if cfg.Watchdog.IntervalSeconds != 60 { + t.Errorf("config.exemple.yml : interval_seconds attendu 60, obtenu %d", cfg.Watchdog.IntervalSeconds) + } + if cfg.QEMU.QMPDir == "" { + t.Error("config.exemple.yml devrait définir qemu.qmp_dir") + } +} diff --git a/internal/config/agent/struct.go b/internal/config/agent/struct.go index 8b5d406..92a8677 100644 --- a/internal/config/agent/struct.go +++ b/internal/config/agent/struct.go @@ -44,6 +44,10 @@ type Config struct { MonitorDir string `mapstructure:"monitor_dir"` QMPDir string `mapstructure:"qmp_dir"` } `mapstructure:"qemu"` + Watchdog struct { + Enabled bool `mapstructure:"enabled"` + IntervalSeconds int `mapstructure:"interval_seconds"` + } `mapstructure:"watchdog"` DefaultInterface string `mapstructure:"default_interface"` Interfaces map[string]string `mapstructure:"interfaces"` } @@ -69,6 +73,8 @@ func LoadConfig(path string) (*Config, error) { v.SetDefault("qemu.serial_dir", "/run/two/vms/serial") v.SetDefault("qemu.monitor_dir", "/run/two/vms/monitor") v.SetDefault("qemu.qmp_dir", "/run/two/vms/qmp") + v.SetDefault("watchdog.enabled", false) + v.SetDefault("watchdog.interval_seconds", 60) v.SetDefault("admin.enabled", false) v.SetDefault("admin.address", "127.0.0.1") v.SetDefault("admin.port", 9091) From d8ece59d1ceb99da7167d65c2d818cfe33599911 Mon Sep 17 00:00:00 2001 From: GnomeZworc Date: Tue, 18 Aug 2026 00:05:04 +0200 Subject: [PATCH 2/4] f-29: watchdog: change error messages #29 Signed-off-by: GnomeZworc --- internal/watchdog/check_subnet.go | 24 ++++++++++++------------ internal/watchdog/check_subnet_test.go | 16 ++++++++-------- internal/watchdog/check_vm.go | 22 +++++++++++----------- internal/watchdog/check_vm_test.go | 12 ++++++------ internal/watchdog/check_vpc.go | 12 ++++++------ internal/watchdog/check_vpc_test.go | 4 ++-- internal/watchdog/notify/stderr.go | 2 +- internal/watchdog/resources.go | 4 ++-- internal/watchdog/watchdog.go | 16 ++++++++-------- internal/watchdog/watchdog_test.go | 14 +++++++------- 10 files changed, 63 insertions(+), 63 deletions(-) diff --git a/internal/watchdog/check_subnet.go b/internal/watchdog/check_subnet.go index 4638791..fafb2ff 100644 --- a/internal/watchdog/check_subnet.go +++ b/internal/watchdog/check_subnet.go @@ -24,7 +24,7 @@ const ( func subnetIfaceNames(subnetName string) (hostVeth, nsVeth, bridge string, err error) { parts := strings.SplitN(subnetName, "-", 2) if len(parts) < 2 || parts[1] == "" { - return "", "", "", fmt.Errorf("nom de subnet %q sans identifiant après le tiret, interfaces indéductibles", subnetName) + return "", "", "", fmt.Errorf("subnet name %q has no identifier after the dash, interface names cannot be derived", subnetName) } id := parts[1] return "v-" + id + "-e", "v-" + id + "-i", "br-" + id, nil @@ -37,13 +37,13 @@ func dnsmasqName(vpc, bridge string) string { func CheckSubnets(db *badger.DB, u unitChecker, n notify.Notifier) error { pairs, err := kv.ListByPrefix(db, prefixSubnet) if err != nil { - return fmt.Errorf("watchdog: lecture des subnets: %w", err) + return fmt.Errorf("watchdog: listing subnets: %w", err) } for _, name := range resourceNames(pairs, prefixSubnet) { st, err := state.Get(db, prefixSubnet+name) if err != nil { - n.Notify(kindSubnet, name, fmt.Sprintf("état illisible en base: %v", err)) + n.Notify(kindSubnet, name, fmt.Sprintf("state unreadable in database: %v", err)) continue } if st != state.Running { @@ -63,13 +63,13 @@ func checkSubnet(db *badger.DB, name string, u unitChecker, n notify.Notifier) { vpc, err := kv.GetFromDB(db, prefixSubnet+name+"/vpc") if err != nil { - n.Notify(kindSubnet, name, fmt.Sprintf("vpc illisible en base: %v", err)) + n.Notify(kindSubnet, name, fmt.Sprintf("vpc unreadable in database: %v", err)) return } mode, err := kv.GetFromDB(db, prefixSubnet+name+"/mode") if err != nil { - n.Notify(kindSubnet, name, fmt.Sprintf("mode illisible en base: %v", err)) + n.Notify(kindSubnet, name, fmt.Sprintf("mode unreadable in database: %v", err)) return } @@ -85,7 +85,7 @@ func checkSubnet(db *badger.DB, name string, u unitChecker, n notify.Notifier) { checkVxlanIface(db, name, n) case modeBridge: default: - n.Notify(kindSubnet, name, fmt.Sprintf("mode inconnu %q", mode)) + n.Notify(kindSubnet, name, fmt.Sprintf("unknown mode %q", mode)) } checkSubnetNetns(name, vpc, nsVeth, bridge, n) @@ -93,7 +93,7 @@ func checkSubnet(db *badger.DB, name string, u unitChecker, n notify.Notifier) { dnsName := dnsmasqName(vpc, bridge) conf := filepath.Join(dhcp.DefaultConfDir, dnsName+".conf") if _, err := os.Stat(conf); err != nil { - n.Notify(kindSubnet, name, fmt.Sprintf("config dnsmasq absente (%s): %v", conf, err)) + n.Notify(kindSubnet, name, fmt.Sprintf("dnsmasq config missing (%s): %v", conf, err)) } checkUnit(kindSubnet, name, "dnsmasq@"+dnsName+".service", u, n) @@ -102,12 +102,12 @@ func checkSubnet(db *badger.DB, name string, u unitChecker, n notify.Notifier) { func checkVxlanIface(db *badger.DB, name string, n notify.Notifier) { raw, err := kv.GetFromDB(db, prefixSubnet+name+"/vxlan_id") if err != nil { - n.Notify(kindSubnet, name, fmt.Sprintf("vxlan_id illisible en base: %v", err)) + n.Notify(kindSubnet, name, fmt.Sprintf("vxlan_id unreadable in database: %v", err)) return } id, err := strconv.Atoi(raw) if err != nil { - n.Notify(kindSubnet, name, fmt.Sprintf("vxlan_id invalide %q: %v", raw, err)) + n.Notify(kindSubnet, name, fmt.Sprintf("invalid vxlan_id %q: %v", raw, err)) return } if p := linkProblem(fmt.Sprintf("vxlan-%d", id)); p != "" { @@ -117,18 +117,18 @@ func checkVxlanIface(db *badger.DB, name string, n notify.Notifier) { func checkSubnetNetns(name, vpc, nsVeth, bridge string, n notify.Notifier) { if !netns.Exist(vpc) { - n.Notify(kindSubnet, name, "netns "+vpc+" absent (/var/run/netns/"+vpc+")") + n.Notify(kindSubnet, name, "netns "+vpc+" missing (/var/run/netns/"+vpc+")") return } if err := netns.Call(vpc, func() error { for _, iface := range []string{nsVeth, bridge} { if p := linkProblem(iface); p != "" { - n.Notify(kindSubnet, name, p+" (dans le netns "+vpc+")") + n.Notify(kindSubnet, name, p+" (in netns "+vpc+")") } } return nil }); err != nil { - n.Notify(kindSubnet, name, fmt.Sprintf("entrée dans le netns %s impossible: %v", vpc, err)) + n.Notify(kindSubnet, name, fmt.Sprintf("cannot enter netns %s: %v", vpc, err)) } } diff --git a/internal/watchdog/check_subnet_test.go b/internal/watchdog/check_subnet_test.go index dbcd19b..6928e6f 100644 --- a/internal/watchdog/check_subnet_test.go +++ b/internal/watchdog/check_subnet_test.go @@ -85,7 +85,7 @@ func TestCheckSubnets_VPCManquantEnBase(t *testing.T) { if len(got) != 1 { t.Fatalf("attendu 1 notification, obtenu %d : %v", len(got), got) } - if !strings.Contains(got[0].problem, "vpc illisible") { + if !strings.Contains(got[0].problem, "vpc unreadable") { t.Errorf("problem = %q, devrait porter sur le vpc", got[0].problem) } } @@ -100,7 +100,7 @@ func TestCheckSubnets_ModeManquantEnBase(t *testing.T) { t.Fatalf("erreur inattendue: %v", err) } - if !r.hasProblemContaining("mode illisible") { + if !r.hasProblemContaining("mode unreadable") { t.Errorf("devrait signaler un mode illisible, obtenu %v", r.calls) } } @@ -114,7 +114,7 @@ func TestCheckSubnets_ModeInconnu(t *testing.T) { t.Fatalf("erreur inattendue: %v", err) } - if !r.hasProblemContaining(`mode inconnu "macvlan"`) { + if !r.hasProblemContaining(`unknown mode "macvlan"`) { t.Errorf("devrait signaler un mode inconnu, obtenu %v", r.calls) } } @@ -147,7 +147,7 @@ func TestCheckSubnets_ModeVxlanSansVxlanID(t *testing.T) { t.Fatalf("erreur inattendue: %v", err) } - if !r.hasProblemContaining("vxlan_id illisible") { + if !r.hasProblemContaining("vxlan_id unreadable") { t.Errorf("devrait signaler un vxlan_id illisible, obtenu %v", r.calls) } } @@ -162,7 +162,7 @@ func TestCheckSubnets_ModeVxlanVxlanIDInvalide(t *testing.T) { t.Fatalf("erreur inattendue: %v", err) } - if !r.hasProblemContaining("vxlan_id invalide") { + if !r.hasProblemContaining("invalid vxlan_id") { t.Errorf("devrait signaler un vxlan_id invalide, obtenu %v", r.calls) } } @@ -191,7 +191,7 @@ func TestCheckSubnets_ConfigDnsmasqAbsente(t *testing.T) { t.Fatalf("erreur inattendue: %v", err) } - if !r.hasProblemContaining("config dnsmasq absente") { + if !r.hasProblemContaining("dnsmasq config missing") { t.Errorf("devrait signaler la config dnsmasq absente, obtenu %v", r.calls) } if !r.hasProblemContaining("vp-admin_br-000042.conf") { @@ -242,7 +242,7 @@ func TestCheckSubnets_UnitIllisible(t *testing.T) { t.Fatalf("erreur inattendue: %v", err) } - if !r.hasProblemContaining("unit dnsmasq@vp-admin_br-000042.service illisible") { + if !r.hasProblemContaining("unit dnsmasq@vp-admin_br-000042.service unreadable") { t.Errorf("devrait signaler l'unit illisible, obtenu %v", r.calls) } } @@ -271,7 +271,7 @@ func TestCheckSubnets_EtatCorrompuNInterrompPasLaBoucle(t *testing.T) { t.Fatalf("un état corrompu ne doit pas faire échouer CheckSubnets: %v", err) } - if !strings.Contains(strings.Join(problems(r.forName("br-corrompu")), " "), "état illisible") { + if !strings.Contains(strings.Join(problems(r.forName("br-corrompu")), " "), "state unreadable") { t.Errorf("devrait signaler l'état corrompu, obtenu %v", r.calls) } if len(r.forName("br-000042")) == 0 { diff --git a/internal/watchdog/check_vm.go b/internal/watchdog/check_vm.go index 8660eb1..3c850a3 100644 --- a/internal/watchdog/check_vm.go +++ b/internal/watchdog/check_vm.go @@ -23,18 +23,18 @@ func tapName(tapID int) string { func CheckVMs(db *badger.DB, cfg *configuration.Config, u unitChecker, n notify.Notifier) error { if cfg == nil { - return errors.New("watchdog: configuration requise pour vérifier les VMs") + return errors.New("watchdog: configuration required to check vms") } pairs, err := kv.ListByPrefix(db, prefixVM) if err != nil { - return fmt.Errorf("watchdog: lecture des vm: %w", err) + return fmt.Errorf("watchdog: listing vms: %w", err) } for _, name := range resourceNames(pairs, prefixVM) { st, err := state.Get(db, prefixVM+name) if err != nil { - n.Notify(kindVM, name, fmt.Sprintf("état illisible en base: %v", err)) + n.Notify(kindVM, name, fmt.Sprintf("state unreadable in database: %v", err)) continue } if st != state.Running { @@ -48,13 +48,13 @@ func CheckVMs(db *badger.DB, cfg *configuration.Config, u unitChecker, n notify. func checkVM(db *badger.DB, cfg *configuration.Config, name string, u unitChecker, n notify.Notifier) { subnetName, err := kv.GetFromDB(db, prefixVM+name+"/subnet") if err != nil { - n.Notify(kindVM, name, fmt.Sprintf("subnet illisible en base: %v", err)) + n.Notify(kindVM, name, fmt.Sprintf("subnet unreadable in database: %v", err)) return } vpc, err := kv.GetFromDB(db, prefixSubnet+subnetName+"/vpc") if err != nil { - n.Notify(kindVM, name, fmt.Sprintf("vpc du subnet %s illisible en base: %v", subnetName, err)) + n.Notify(kindVM, name, fmt.Sprintf("vpc of subnet %s unreadable in database: %v", subnetName, err)) return } @@ -67,33 +67,33 @@ func checkVM(db *badger.DB, cfg *configuration.Config, name string, u unitChecke func checkVMTap(db *badger.DB, name, vpc string, n notify.Notifier) { raw, err := kv.GetFromDB(db, prefixVM+name+"/tap_id") if err != nil { - n.Notify(kindVM, name, fmt.Sprintf("tap_id illisible en base: %v", err)) + n.Notify(kindVM, name, fmt.Sprintf("tap_id unreadable in database: %v", err)) return } tapID, err := strconv.Atoi(raw) if err != nil { - n.Notify(kindVM, name, fmt.Sprintf("tap_id invalide %q: %v", raw, err)) + n.Notify(kindVM, name, fmt.Sprintf("invalid tap_id %q: %v", raw, err)) return } if !netns.Exist(vpc) { - n.Notify(kindVM, name, "netns "+vpc+" absent (/var/run/netns/"+vpc+")") + n.Notify(kindVM, name, "netns "+vpc+" missing (/var/run/netns/"+vpc+")") return } if err := netns.Call(vpc, func() error { if p := linkProblem(tapName(tapID)); p != "" { - n.Notify(kindVM, name, p+" (dans le netns "+vpc+")") + n.Notify(kindVM, name, p+" (in netns "+vpc+")") } return nil }); err != nil { - n.Notify(kindVM, name, fmt.Sprintf("entrée dans le netns %s impossible: %v", vpc, err)) + n.Notify(kindVM, name, fmt.Sprintf("cannot enter netns %s: %v", vpc, err)) } } func checkVMQemu(cfg *configuration.Config, name string, n notify.Notifier) { sock := filepath.Join(cfg.QEMU.QMPDir, name+".sock") if _, err := qmp.Send(sock, nil); err != nil { - n.Notify(kindVM, name, fmt.Sprintf("qemu ne répond pas sur %s: %v", sock, err)) + n.Notify(kindVM, name, fmt.Sprintf("qemu not responding on %s: %v", sock, err)) } } diff --git a/internal/watchdog/check_vm_test.go b/internal/watchdog/check_vm_test.go index 8bf24a4..eb62b59 100644 --- a/internal/watchdog/check_vm_test.go +++ b/internal/watchdog/check_vm_test.go @@ -86,7 +86,7 @@ func TestCheckVMs_SubnetManquantEnBase(t *testing.T) { if len(got) != 1 { t.Fatalf("attendu 1 notification, obtenu %d : %v", len(got), got) } - if !strings.Contains(got[0].problem, "subnet illisible") { + if !strings.Contains(got[0].problem, "subnet unreadable") { t.Errorf("problem = %q, devrait porter sur le subnet", got[0].problem) } } @@ -101,7 +101,7 @@ func TestCheckVMs_VPCDuSubnetManquant(t *testing.T) { t.Fatalf("erreur inattendue: %v", err) } - if !r.hasProblemContaining("vpc du subnet br-000042 illisible") { + if !r.hasProblemContaining("vpc of subnet br-000042 unreadable") { t.Errorf("devrait signaler le vpc introuvable, obtenu %v", r.calls) } } @@ -115,7 +115,7 @@ func TestCheckVMs_TapIDManquant(t *testing.T) { t.Fatalf("erreur inattendue: %v", err) } - if !r.hasProblemContaining("tap_id illisible") { + if !r.hasProblemContaining("tap_id unreadable") { t.Errorf("devrait signaler un tap_id illisible, obtenu %v", r.calls) } } @@ -129,7 +129,7 @@ func TestCheckVMs_TapIDInvalide(t *testing.T) { t.Fatalf("erreur inattendue: %v", err) } - if !r.hasProblemContaining("tap_id invalide") { + if !r.hasProblemContaining("invalid tap_id") { t.Errorf("devrait signaler un tap_id invalide, obtenu %v", r.calls) } } @@ -144,7 +144,7 @@ func TestCheckVMs_QemuNeRepondPas(t *testing.T) { t.Fatalf("erreur inattendue: %v", err) } - if !r.hasProblemContaining("qemu ne répond pas") { + if !r.hasProblemContaining("qemu not responding") { t.Errorf("devrait signaler que qemu ne répond pas, obtenu %v", r.calls) } if !r.hasProblemContaining("i-test1.sock") { @@ -220,7 +220,7 @@ func TestCheckVMs_EtatCorrompuNInterrompPasLaBoucle(t *testing.T) { t.Fatalf("un état corrompu ne doit pas faire échouer CheckVMs: %v", err) } - if !strings.Contains(strings.Join(problems(r.forName("i-corrompu")), " "), "état illisible") { + if !strings.Contains(strings.Join(problems(r.forName("i-corrompu")), " "), "state unreadable") { t.Errorf("devrait signaler l'état corrompu, obtenu %v", r.calls) } if len(r.forName("i-test1")) == 0 { diff --git a/internal/watchdog/check_vpc.go b/internal/watchdog/check_vpc.go index 5cf252d..c344561 100644 --- a/internal/watchdog/check_vpc.go +++ b/internal/watchdog/check_vpc.go @@ -17,7 +17,7 @@ const vpcBridge = "br-public" func vpcIfaceNames(vpcName string) (host, ns string, err error) { parts := strings.SplitN(vpcName, "-", 2) if len(parts) < 2 || parts[1] == "" { - return "", "", fmt.Errorf("nom de VPC %q sans identifiant après le tiret, interfaces indéductibles", vpcName) + return "", "", fmt.Errorf("vpc name %q has no identifier after the dash, interface names cannot be derived", vpcName) } return "vp-" + parts[1] + "-e", "vp-" + parts[1] + "-i", nil } @@ -25,13 +25,13 @@ func vpcIfaceNames(vpcName string) (host, ns string, err error) { func CheckVPCs(db *badger.DB, n notify.Notifier) error { pairs, err := kv.ListByPrefix(db, prefixVPC) if err != nil { - return fmt.Errorf("watchdog: lecture des vpc: %w", err) + return fmt.Errorf("watchdog: listing vpcs: %w", err) } for _, name := range resourceNames(pairs, prefixVPC) { st, err := state.Get(db, prefixVPC+name) if err != nil { - n.Notify(kindVPC, name, fmt.Sprintf("état illisible en base: %v", err)) + n.Notify(kindVPC, name, fmt.Sprintf("state unreadable in database: %v", err)) continue } if st != state.Running { @@ -50,7 +50,7 @@ func checkVPC(name string, n notify.Notifier) { } if !netns.Exist(name) { - n.Notify(kindVPC, name, "netns absent (/var/run/netns/"+name+")") + n.Notify(kindVPC, name, "netns missing (/var/run/netns/"+name+")") return } @@ -61,11 +61,11 @@ func checkVPC(name string, n notify.Notifier) { if err := netns.Call(name, func() error { for _, iface := range []string{nsVeth, vpcBridge} { if p := linkProblem(iface); p != "" { - n.Notify(kindVPC, name, p+" (dans le netns)") + n.Notify(kindVPC, name, p+" (in netns)") } } return nil }); err != nil { - n.Notify(kindVPC, name, fmt.Sprintf("entrée dans le netns impossible: %v", err)) + n.Notify(kindVPC, name, fmt.Sprintf("cannot enter netns: %v", err)) } } diff --git a/internal/watchdog/check_vpc_test.go b/internal/watchdog/check_vpc_test.go index 1b37248..3ae8a07 100644 --- a/internal/watchdog/check_vpc_test.go +++ b/internal/watchdog/check_vpc_test.go @@ -122,7 +122,7 @@ func TestCheckVPCs_EtatCorrompu(t *testing.T) { if len(got) != 1 { t.Fatalf("attendu 1 notification pour l'état corrompu, obtenu %d", len(got)) } - if !strings.Contains(got[0].problem, "état illisible") { + if !strings.Contains(got[0].problem, "state unreadable") { t.Errorf("problem = %q, devrait mentionner un état illisible", got[0].problem) } if len(r.forName("vp-suivant")) == 0 { @@ -143,7 +143,7 @@ func TestCheckVPCs_NomIndeductibleNePaniquePas(t *testing.T) { if len(got) != 1 { t.Fatalf("attendu 1 notification, obtenu %d : %v", len(got), got) } - if !strings.Contains(got[0].problem, "indéductibles") { + if !strings.Contains(got[0].problem, "cannot be derived") { t.Errorf("problem = %q, devrait porter sur les interfaces indéductibles", got[0].problem) } } diff --git a/internal/watchdog/notify/stderr.go b/internal/watchdog/notify/stderr.go index 5436b2c..cac8174 100644 --- a/internal/watchdog/notify/stderr.go +++ b/internal/watchdog/notify/stderr.go @@ -16,7 +16,7 @@ func NewStderr(l *slog.Logger) *StderrNotifier { } func (n *StderrNotifier) Notify(kind, name, problem string) { - n.logger.Error("watchdog: incohérence détectée", + n.logger.Error("watchdog: inconsistency detected", "kind", kind, "name", name, "problem", problem, diff --git a/internal/watchdog/resources.go b/internal/watchdog/resources.go index d5e3fd4..37219a5 100644 --- a/internal/watchdog/resources.go +++ b/internal/watchdog/resources.go @@ -52,7 +52,7 @@ func checkUnit(kind, name, unit string, u unitChecker, n notify.Notifier) { } st, err := u.Status(unit) if err != nil { - n.Notify(kind, name, fmt.Sprintf("unit %s illisible: %v", unit, err)) + n.Notify(kind, name, fmt.Sprintf("unit %s unreadable: %v", unit, err)) return } if st.ActiveState != "active" { @@ -64,7 +64,7 @@ func linkProblem(iface string) string { up, err := netif.LinkIsUp(iface) switch { case err != nil: - return fmt.Sprintf("interface %s introuvable: %v", iface, err) + return fmt.Sprintf("interface %s not found: %v", iface, err) case !up: return fmt.Sprintf("interface %s down", iface) } diff --git a/internal/watchdog/watchdog.go b/internal/watchdog/watchdog.go index 0dbb85b..3ab1d52 100644 --- a/internal/watchdog/watchdog.go +++ b/internal/watchdog/watchdog.go @@ -28,7 +28,7 @@ func New(db *badger.DB, cfg *configuration.Config, n notify.Notifier, logger *sl logger = slog.Default() } if interval <= 0 { - logger.Warn("watchdog: intervalle invalide, valeur par défaut appliquée", + logger.Warn("watchdog: invalid interval, default applied", "interval", interval, "default", defaultInterval) interval = defaultInterval } @@ -45,12 +45,12 @@ func (w *Watchdog) Run(ctx context.Context) { ticker := time.NewTicker(w.interval) defer ticker.Stop() - w.logger.Info("watchdog: démarrage", "interval", w.interval) + w.logger.Info("watchdog: starting", "interval", w.interval) for { select { case <-ctx.Done(): - w.logger.Info("watchdog: arrêt") + w.logger.Info("watchdog: stopping") return case <-ticker.C: w.tick() @@ -63,13 +63,13 @@ func (w *Watchdog) tick() { defer closeUnits() if err := CheckVPCs(w.db, w.notifier); err != nil { - w.logger.Error("watchdog: vérification des vpc", "err", err) + w.logger.Error("watchdog: vpc check failed", "err", err) } if err := CheckSubnets(w.db, u, w.notifier); err != nil { - w.logger.Error("watchdog: vérification des subnets", "err", err) + w.logger.Error("watchdog: subnet check failed", "err", err) } if err := CheckVMs(w.db, w.cfg, u, w.notifier); err != nil { - w.logger.Error("watchdog: vérification des vm", "err", err) + w.logger.Error("watchdog: vm check failed", "err", err) } } @@ -77,13 +77,13 @@ func (w *Watchdog) units() (unitChecker, func()) { m, err := systemd.New() if err != nil { if !w.dbusDown { - w.logger.Warn("watchdog: connexion systemd impossible, vérification des units désactivée", "err", err) + w.logger.Warn("watchdog: systemd unreachable, unit checks disabled", "err", err) w.dbusDown = true } return nil, func() {} } if w.dbusDown { - w.logger.Info("watchdog: connexion systemd rétablie") + w.logger.Info("watchdog: systemd connection restored") w.dbusDown = false } return m, m.Close diff --git a/internal/watchdog/watchdog_test.go b/internal/watchdog/watchdog_test.go index c8cea47..019ddb5 100644 --- a/internal/watchdog/watchdog_test.go +++ b/internal/watchdog/watchdog_test.go @@ -61,7 +61,7 @@ func TestNew_IntervalleInvalideUtiliseLeDefaut(t *testing.T) { } } - if !strings.Contains(buf.String(), "intervalle invalide") { + if !strings.Contains(buf.String(), "invalid interval") { t.Error("un intervalle invalide devrait être signalé dans les logs") } } @@ -143,11 +143,11 @@ func TestRun_ErreurDeBaseLogueeEtBoucleContinue(t *testing.T) { for _, line := range logLines(t, buf) { msg, _ := line["msg"].(string) switch { - case strings.Contains(msg, "vérification des vpc"): + case strings.Contains(msg, "vpc check failed"): vpc = true - case strings.Contains(msg, "vérification des subnets"): + case strings.Contains(msg, "subnet check failed"): subnet = true - case strings.Contains(msg, "vérification des vm"): + case strings.Contains(msg, "vm check failed"): vm = true } } @@ -168,10 +168,10 @@ func TestRun_LogueDemarrageEtArret(t *testing.T) { w.Run(ctx) out := buf.String() - if !strings.Contains(out, "watchdog: démarrage") { + if !strings.Contains(out, "watchdog: starting") { t.Error("le démarrage devrait être logué") } - if !strings.Contains(out, "watchdog: arrêt") { + if !strings.Contains(out, "watchdog: stopping") { t.Error("l'arrêt devrait être logué") } } @@ -190,7 +190,7 @@ func TestUnits_ConnexionSystemdIndisponibleSignaleeUneSeuleFois(t *testing.T) { var warnings int for _, line := range logLines(t, buf) { - if msg, _ := line["msg"].(string); strings.Contains(msg, "connexion systemd impossible") { + if msg, _ := line["msg"].(string); strings.Contains(msg, "systemd unreachable") { warnings++ } } From 1ce84ffd58516831c2102cba1aaff7ee0c7b5e52 Mon Sep 17 00:00:00 2001 From: GnomeZworc Date: Tue, 18 Aug 2026 00:22:57 +0200 Subject: [PATCH 3/4] f-29: watchdog: call this one #29 Signed-off-by: GnomeZworc --- cmd/agent/main.go | 10 ++++++++++ 1 file changed, 10 insertions(+) diff --git a/cmd/agent/main.go b/cmd/agent/main.go index 0425d4e..cd14678 100644 --- a/cmd/agent/main.go +++ b/cmd/agent/main.go @@ -1,15 +1,19 @@ package main import ( + "context" "flag" "fmt" "log/slog" + "time" agentapi "git.g3e.fr/syonad/two/internal/api/agent" configuration "git.g3e.fr/syonad/two/internal/config/agent" dispatcher "git.g3e.fr/syonad/two/internal/dispatcher/agent" "git.g3e.fr/syonad/two/internal/migration" agentmetrics "git.g3e.fr/syonad/two/internal/prometheus/agent" + "git.g3e.fr/syonad/two/internal/watchdog" + "git.g3e.fr/syonad/two/internal/watchdog/notify" "git.g3e.fr/syonad/two/pkg/db/kv" "git.g3e.fr/syonad/two/pkg/logger" promserver "git.g3e.fr/syonad/two/pkg/prometheus" @@ -63,6 +67,12 @@ func main() { adminAddr := fmt.Sprintf("%s:%d", cfg.Admin.Address, cfg.Admin.Port) go kv.NewAdminServer(db, log.With(slog.String("component", "admin"))).Start(adminAddr) } + if cfg.Watchdog.Enabled { + wlog := log.With(slog.String("component", "watchdog")) + go watchdog.New(db, cfg, notify.NewStderr(wlog), wlog, + time.Duration(cfg.Watchdog.IntervalSeconds)*time.Second, + ).Run(context.Background()) + } select {} } From cdbb5012b11dbcd0b97dc80fafcdcf574382b2fc Mon Sep 17 00:00:00 2001 From: GnomeZworc Date: Tue, 18 Aug 2026 00:41:01 +0200 Subject: [PATCH 4/4] f-29: stop: add a proper stop Signed-off-by: GnomeZworc --- cmd/agent/main.go | 72 ++++++++++++++++++++-- cmd/agent/main_test.go | 98 ++++++++++++++++++++++++++++++ internal/api/agent/helpers_test.go | 2 +- internal/api/agent/server.go | 23 ++++--- pkg/db/kv/admin_server.go | 23 ++++--- pkg/prometheus/server.go | 29 +++++++-- pkg/worker/queue.go | 38 +++++++++++- pkg/worker/queue_test.go | 76 +++++++++++++++++++++++ systemd/agent.service | 1 + 9 files changed, 333 insertions(+), 29 deletions(-) create mode 100644 cmd/agent/main_test.go diff --git a/cmd/agent/main.go b/cmd/agent/main.go index cd14678..1cea7ec 100644 --- a/cmd/agent/main.go +++ b/cmd/agent/main.go @@ -5,6 +5,9 @@ import ( "flag" "fmt" "log/slog" + "os" + "os/signal" + "syscall" "time" agentapi "git.g3e.fr/syonad/two/internal/api/agent" @@ -21,6 +24,8 @@ import ( "github.com/prometheus/client_golang/prometheus" ) +const shutdownTimeout = 20 * time.Second + func main() { confFile := flag.String("config", "/etc/two/agent.yml", "config file path") flag.Parse() @@ -34,7 +39,12 @@ func main() { log := logger.New(cfg.Logger.Level, cfg.Logger.Debug) 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 // d'états, et des ressources transitoires orphelines d'un arrêt précédent. @@ -60,19 +70,69 @@ func main() { "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"))) - 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 { 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 { wlog := log.With(slog.String("component", "watchdog")) go watchdog.New(db, cfg, notify.NewStderr(wlog), wlog, 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 + } } diff --git a/cmd/agent/main_test.go b/cmd/agent/main_test.go new file mode 100644 index 0000000..12291c7 --- /dev/null +++ b/cmd/agent/main_test.go @@ -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") + } +} diff --git a/internal/api/agent/helpers_test.go b/internal/api/agent/helpers_test.go index 206874c..c7162c7 100644 --- a/internal/api/agent/helpers_test.go +++ b/internal/api/agent/helpers_test.go @@ -23,5 +23,5 @@ func newTestServer(t *testing.T) (*Server, *badger.DB) { cfg := &configuration.Config{DefaultInterface: "br-test"} logger := slog.New(slog.NewTextHandler(io.Discard, nil)) d := dispatcher.New(q, db, cfg, logger) - return New(d, db, logger), db + return New(d, db, logger, "127.0.0.1:0"), db } diff --git a/internal/api/agent/server.go b/internal/api/agent/server.go index 4a0fba4..2b892bd 100644 --- a/internal/api/agent/server.go +++ b/internal/api/agent/server.go @@ -1,8 +1,10 @@ package agentapi import ( + "context" "crypto/rand" "encoding/hex" + "errors" "log/slog" "net/http" "time" @@ -15,13 +17,11 @@ type Server struct { dispatcher *dispatcher.Dispatcher db *badger.DB logger *slog.Logger + srv *http.Server } -func New(d *dispatcher.Dispatcher, db *badger.DB, logger *slog.Logger) *Server { - return &Server{dispatcher: d, db: db, logger: logger} -} - -func (s *Server) Start(address string) { +func New(d *dispatcher.Dispatcher, db *badger.DB, logger *slog.Logger, address string) *Server { + s := &Server{dispatcher: d, db: db, logger: logger} mux := http.NewServeMux() mux.HandleFunc("/vpcs", s.VpcsHandler) mux.HandleFunc("/vpcs/", s.VpcByNameHandler) @@ -29,12 +29,21 @@ func (s *Server) Start(address string) { mux.HandleFunc("/subnets/", s.SubnetByNameHandler) mux.HandleFunc("/vms", s.VmsHandler) mux.HandleFunc("/vms/", s.VmByNameHandler) - s.logger.Info("API server listening", "address", address) - if err := http.ListenAndServe(address, s.logMiddleware(mux)); err != nil { + s.srv = &http.Server{Addr: address, Handler: s.logMiddleware(mux)} + 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) } } +func (s *Server) Shutdown(ctx context.Context) error { + return s.srv.Shutdown(ctx) +} + type statusWriter struct { http.ResponseWriter status int diff --git a/pkg/db/kv/admin_server.go b/pkg/db/kv/admin_server.go index 1216244..db37f6c 100644 --- a/pkg/db/kv/admin_server.go +++ b/pkg/db/kv/admin_server.go @@ -1,6 +1,8 @@ package kv import ( + "context" + "errors" "fmt" "log/slog" "net/http" @@ -12,21 +14,28 @@ import ( type AdminServer struct { db *badger.DB logger *slog.Logger + srv *http.Server } -func NewAdminServer(db *badger.DB, logger *slog.Logger) *AdminServer { - return &AdminServer{db: db, logger: logger} -} - -func (s *AdminServer) Start(address string) { +func NewAdminServer(db *badger.DB, logger *slog.Logger, address string) *AdminServer { + s := &AdminServer{db: db, logger: logger} mux := http.NewServeMux() mux.HandleFunc("/db", s.dbHandler) - s.logger.Info("admin server listening", "address", address) - if err := http.ListenAndServe(address, mux); err != nil { + s.srv = &http.Server{Addr: address, Handler: mux} + 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) } } +func (s *AdminServer) Shutdown(ctx context.Context) error { + return s.srv.Shutdown(ctx) +} + func (s *AdminServer) dbHandler(w http.ResponseWriter, r *http.Request) { if r.Method != http.MethodGet { http.Error(w, "method not allowed", http.StatusMethodNotAllowed) diff --git a/pkg/prometheus/server.go b/pkg/prometheus/server.go index 9f8e557..f6c34b3 100644 --- a/pkg/prometheus/server.go +++ b/pkg/prometheus/server.go @@ -1,6 +1,8 @@ package promserver import ( + "context" + "errors" "log" "net/http" @@ -8,13 +10,30 @@ import ( "github.com/prometheus/client_golang/prometheus/promhttp" ) -// Start launches the Prometheus metrics HTTP server on the given address. -// The provided registry is used to expose metrics at /metrics. -func Start(address string, registry *prometheus.Registry) { +// Server exposes the Prometheus metrics endpoint. +type Server struct { + 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.Handle("/metrics", promhttp.HandlerFor(registry, promhttp.HandlerOpts{ EnableOpenMetrics: true, })) - log.Printf("Prometheus server listening on %s", address) - log.Fatal(http.ListenAndServe(address, mux)) + return &Server{srv: &http.Server{Addr: address, Handler: 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) } diff --git a/pkg/worker/queue.go b/pkg/worker/queue.go index 109c726..40e1443 100644 --- a/pkg/worker/queue.go +++ b/pkg/worker/queue.go @@ -1,13 +1,19 @@ package worker -import "log" +import ( + "log" + "sync" +) // 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 + tasks chan Task + wg sync.WaitGroup + mu sync.RWMutex + stopped bool } // 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)} } -// 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) { + q.mu.RLock() + defer q.mu.RUnlock() + + if q.stopped { + log.Print("worker: queue stopped, task rejected") + return + } q.tasks <- t } @@ -24,10 +38,28 @@ func (q *Queue) Submit(t Task) { func (q *Queue) Start(n int) { log.Printf("worker: starting %d workers", n) for i := range n { + q.wg.Add(1) go func(id int) { + defer q.wg.Done() for task := range q.tasks { task() } }(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() +} diff --git a/pkg/worker/queue_test.go b/pkg/worker/queue_test.go index 5353b06..4e463d6 100644 --- a/pkg/worker/queue_test.go +++ b/pkg/worker/queue_test.go @@ -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") } } + +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() +} diff --git a/systemd/agent.service b/systemd/agent.service index 37715c4..95ccf7e 100644 --- a/systemd/agent.service +++ b/systemd/agent.service @@ -5,6 +5,7 @@ After=network.target [Service] Type=simple ExecStart=/opt/two/bin/agent +TimeoutStopSec=30 [Install] WantedBy=multi-user.target