From 69c7a352f64a3d1993f32b13fc1d2dc77fc77e83 Mon Sep 17 00:00:00 2001 From: mrhid6 Date: Tue, 21 Jul 2026 14:19:59 +0100 Subject: [PATCH] feat(agent): agent-run monitor scheduler --- agent/internal/checker/checker.go | 222 ++++++++++++++++++++++++++++ agent/internal/monitors/monitors.go | 165 +++++++++++++++++++++ agent/internal/sync/sync.go | 4 + 3 files changed, 391 insertions(+) create mode 100644 agent/internal/checker/checker.go create mode 100644 agent/internal/monitors/monitors.go diff --git a/agent/internal/checker/checker.go b/agent/internal/checker/checker.go new file mode 100644 index 0000000..c4d49dd --- /dev/null +++ b/agent/internal/checker/checker.go @@ -0,0 +1,222 @@ +// Package checker runs service checks (http/tcp/icmp/tls) and returns a uniform +// Result. It has no dependency on models or pb so it can be duplicated verbatim +// into the agent module (agent-run monitors) — callers map their own monitor +// representation onto Spec. +package checker + +import ( + "context" + "crypto/tls" + "fmt" + "io" + "net" + "net/http" + "os" + "strings" + "time" +) + +// Check types (mirror models.Monitor* constants). +const ( + TypeHTTP = "http" + TypeTCP = "tcp" + TypeICMP = "icmp" + TypeTLS = "tls" +) + +// Spec is a self-contained description of a single check. +type Spec struct { + Type string + URL string + Host string + Port int + Method string + ExpectedStatus int + Keyword string + TLSWarnDays int + TimeoutSec int +} + +// Result is the uniform outcome of running a check. +type Result struct { + Up bool + LatencyMs int + Message string + CertExpiry *time.Time +} + +func (s Spec) timeout() time.Duration { + t := s.TimeoutSec + if t <= 0 || t > 10 { + t = 10 + } + return time.Duration(t) * time.Second +} + +// Run executes the check described by s. +func Run(ctx context.Context, s Spec) Result { + switch s.Type { + case TypeHTTP: + return runHTTP(ctx, s) + case TypeTCP: + return runTCP(ctx, s) + case TypeICMP: + return runICMP(ctx, s) + case TypeTLS: + return runTLS(ctx, s) + default: + return Result{Message: "unknown check type: " + s.Type} + } +} + +func runHTTP(ctx context.Context, s Spec) Result { + method := s.Method + if method == "" { + method = http.MethodGet + } + expect := s.ExpectedStatus + if expect == 0 { + expect = 200 + } + client := &http.Client{Timeout: s.timeout()} + start := time.Now() + req, err := http.NewRequestWithContext(ctx, method, s.URL, nil) + if err != nil { + return Result{Message: err.Error()} + } + resp, err := client.Do(req) + if err != nil { + return Result{LatencyMs: msSince(start), Message: err.Error()} + } + defer resp.Body.Close() + res := Result{LatencyMs: msSince(start), Up: true} + if resp.TLS != nil && len(resp.TLS.PeerCertificates) > 0 { + exp := resp.TLS.PeerCertificates[0].NotAfter + res.CertExpiry = &exp + } + if resp.StatusCode != expect { + return Result{LatencyMs: res.LatencyMs, CertExpiry: res.CertExpiry, Message: fmt.Sprintf("status %d (want %d)", resp.StatusCode, expect)} + } + if s.Keyword != "" { + body, _ := io.ReadAll(io.LimitReader(resp.Body, 1<<20)) + if !strings.Contains(string(body), s.Keyword) { + return Result{LatencyMs: res.LatencyMs, CertExpiry: res.CertExpiry, Message: "keyword not found"} + } + } + return res +} + +func runTCP(ctx context.Context, s Spec) Result { + addr := net.JoinHostPort(s.Host, fmt.Sprint(s.Port)) + start := time.Now() + d := net.Dialer{Timeout: s.timeout()} + conn, err := d.DialContext(ctx, "tcp", addr) + if err != nil { + return Result{LatencyMs: msSince(start), Message: err.Error()} + } + conn.Close() + return Result{Up: true, LatencyMs: msSince(start)} +} + +func runTLS(ctx context.Context, s Spec) Result { + port := s.Port + if port == 0 { + port = 443 + } + addr := net.JoinHostPort(s.Host, fmt.Sprint(port)) + start := time.Now() + d := net.Dialer{Timeout: s.timeout()} + conn, err := tls.DialWithDialer(&d, "tcp", addr, &tls.Config{ServerName: s.Host}) + if err != nil { + return Result{LatencyMs: msSince(start), Message: err.Error()} + } + defer conn.Close() + certs := conn.ConnectionState().PeerCertificates + if len(certs) == 0 { + return Result{LatencyMs: msSince(start), Message: "no peer certificate"} + } + exp := certs[0].NotAfter + res := Result{LatencyMs: msSince(start), CertExpiry: &exp} + warn := s.TLSWarnDays + if warn <= 0 { + warn = 14 + } + remaining := time.Until(exp) + if remaining <= 0 { + res.Message = "certificate expired" + return res + } + if remaining <= time.Duration(warn)*24*time.Hour { + res.Message = fmt.Sprintf("certificate expires in %d days", int(remaining.Hours()/24)) + return res + } + res.Up = true + return res +} + +func msSince(t time.Time) int { return int(time.Since(t).Milliseconds()) } + +// runICMP sends a single ICMP echo request and waits for the reply. Requires +// raw-socket privileges (the agent and server run as root). Returns down with a +// descriptive message when the socket cannot be opened or no reply arrives. +func runICMP(ctx context.Context, s Spec) Result { + dst, err := net.ResolveIPAddr("ip4", s.Host) + if err != nil { + return Result{Message: err.Error()} + } + conn, err := net.ListenPacket("ip4:icmp", "0.0.0.0") + if err != nil { + return Result{Message: "icmp socket: " + err.Error()} + } + defer conn.Close() + + id := os.Getpid() & 0xffff + pkt := icmpEcho(id, 1) + deadline := time.Now().Add(s.timeout()) + if d, ok := ctx.Deadline(); ok && d.Before(deadline) { + deadline = d + } + _ = conn.SetDeadline(deadline) + + start := time.Now() + if _, err := conn.WriteTo(pkt, dst); err != nil { + return Result{Message: err.Error()} + } + reply := make([]byte, 1500) + for { + n, peer, err := conn.ReadFrom(reply) + if err != nil { + return Result{LatencyMs: msSince(start), Message: "no reply"} + } + // Skip the IPv4 header (20 bytes) to reach the ICMP message. + if n < 28 || peer.String() != dst.String() { + continue + } + if reply[20] == 0 { // ICMP echo reply type + return Result{Up: true, LatencyMs: msSince(start)} + } + } +} + +func icmpEcho(id, seq int) []byte { + // Type(8)=echo request, Code=0, Checksum, ID, Seq, no payload. + b := []byte{8, 0, 0, 0, byte(id >> 8), byte(id), byte(seq >> 8), byte(seq)} + cs := icmpChecksum(b) + b[2] = byte(cs >> 8) + b[3] = byte(cs) + return b +} + +func icmpChecksum(b []byte) uint16 { + var sum uint32 + for i := 0; i < len(b)-1; i += 2 { + sum += uint32(b[i])<<8 | uint32(b[i+1]) + } + if len(b)%2 == 1 { + sum += uint32(b[len(b)-1]) << 8 + } + for sum>>16 != 0 { + sum = (sum & 0xffff) + (sum >> 16) + } + return ^uint16(sum) +} diff --git a/agent/internal/monitors/monitors.go b/agent/internal/monitors/monitors.go new file mode 100644 index 0000000..69f76a4 --- /dev/null +++ b/agent/internal/monitors/monitors.go @@ -0,0 +1,165 @@ +// Package monitors runs agent-side service checks. It polls the server for the +// monitors assigned to this agent (SyncMonitors), runs each on its own interval +// using the local checker package, and reports results back (ReportChecks). +package monitors + +import ( + "context" + "log" + "sync" + "time" + + "github.com/mrhid6/vantage/agent/internal/checker" + "github.com/mrhid6/vantage/agent/internal/config" + grpcclient "github.com/mrhid6/vantage/agent/internal/grpc" + "github.com/mrhid6/vantage/agent/internal/grpc/pb" +) + +// syncInterval controls how often the agent re-fetches its assigned monitors. +const syncInterval = 30 * time.Second + +type runner struct { + intervalSec int + cancel context.CancelFunc +} + +// Run starts the agent monitor loop and blocks until ctx is cancelled. +func Run(ctx context.Context, cfg *config.Config) { + active := map[string]*runner{} + var mu sync.Mutex + + // results is a shared channel every check writes to; a single reporter + // goroutine batches and ships them so we make one ReportChecks call per tick. + results := make(chan pb.CheckResult, 64) + go reporter(ctx, cfg, results) + + syncOnce := func() { + specs, err := fetchSpecs(cfg) + if err != nil { + log.Printf("monitors: sync: %v", err) + return + } + want := map[string]pb.MonitorSpec{} + for _, s := range specs { + want[s.MonitorId] = s + } + + mu.Lock() + defer mu.Unlock() + for id, r := range active { + s, ok := want[id] + if !ok || s.IntervalSec != r.intervalSec { + r.cancel() + delete(active, id) + } + } + for id, s := range want { + if _, ok := active[id]; ok { + continue + } + rctx, cancel := context.WithCancel(ctx) + active[id] = &runner{intervalSec: s.IntervalSec, cancel: cancel} + go runSpec(rctx, s, results) + } + } + + syncOnce() + t := time.NewTicker(syncInterval) + defer t.Stop() + for { + select { + case <-ctx.Done(): + return + case <-t.C: + syncOnce() + } + } +} + +func fetchSpecs(cfg *config.Config) ([]pb.MonitorSpec, error) { + client, err := grpcclient.New(cfg.ServerURL, cfg.TLS) + if err != nil { + return nil, err + } + defer client.Close() + return client.SyncMonitors(cfg.ServerID, cfg.AgentToken) +} + +func runSpec(ctx context.Context, s pb.MonitorSpec, out chan<- pb.CheckResult) { + interval := time.Duration(s.IntervalSec) * time.Second + if interval <= 0 { + interval = 60 * time.Second + } + spec := checker.Spec{ + Type: s.Type, + URL: s.URL, + Host: s.Host, + Port: s.Port, + Method: s.Method, + ExpectedStatus: s.ExpectedStatus, + Keyword: s.Keyword, + TLSWarnDays: s.TLSWarnDays, + TimeoutSec: s.IntervalSec, + } + + run := func() { + res := checker.Run(ctx, spec) + cr := pb.CheckResult{MonitorId: s.MonitorId, Up: res.Up, LatencyMs: res.LatencyMs, Message: res.Message} + if res.CertExpiry != nil { + cr.CertExpiryUnix = res.CertExpiry.Unix() + } + select { + case out <- cr: + case <-ctx.Done(): + } + } + + run() + t := time.NewTicker(interval) + defer t.Stop() + for { + select { + case <-ctx.Done(): + return + case <-t.C: + run() + } + } +} + +// reporter batches results on a short interval and ships each batch in one call. +func reporter(ctx context.Context, cfg *config.Config, in <-chan pb.CheckResult) { + t := time.NewTicker(5 * time.Second) + defer t.Stop() + var batch []pb.CheckResult + flush := func() { + if len(batch) == 0 { + return + } + client, err := grpcclient.New(cfg.ServerURL, cfg.TLS) + if err != nil { + log.Printf("monitors: report dial: %v", err) + batch = nil + return + } + if err := client.ReportChecks(cfg.ServerID, cfg.AgentToken, batch); err != nil { + log.Printf("monitors: report: %v", err) + } + client.Close() + batch = nil + } + for { + select { + case <-ctx.Done(): + flush() + return + case r := <-in: + batch = append(batch, r) + if len(batch) >= 32 { + flush() + } + case <-t.C: + flush() + } + } +} diff --git a/agent/internal/sync/sync.go b/agent/internal/sync/sync.go index 5d56f58..9e69169 100644 --- a/agent/internal/sync/sync.go +++ b/agent/internal/sync/sync.go @@ -23,6 +23,7 @@ import ( "github.com/mrhid6/vantage/agent/internal/grpc/pb" "github.com/mrhid6/vantage/agent/internal/inventory" "github.com/mrhid6/vantage/agent/internal/keys" + "github.com/mrhid6/vantage/agent/internal/monitors" "github.com/mrhid6/vantage/agent/internal/updates" ) @@ -72,6 +73,9 @@ func Run(ctx context.Context, cfg *config.Config, version string) error { // Report host inventory: metrics every 30s, full static snapshot every 15 min. go runInventory(ctx, cfg) + // Run agent-side service monitors assigned to this server. + go monitors.Run(ctx, cfg) + ticker := time.NewTicker(cfg.PollInterval) defer ticker.Stop()