From e019493087a7dc73bdc8891c7b320f9b62fbc336 Mon Sep 17 00:00:00 2001 From: mrhid6 Date: Tue, 21 Jul 2026 14:10:46 +0100 Subject: [PATCH] feat(server): monitor model + checker package --- server/internal/checker/checker.go | 222 +++++++++++++++++++++++++++++ server/internal/models/monitor.go | 77 ++++++++++ 2 files changed, 299 insertions(+) create mode 100644 server/internal/checker/checker.go create mode 100644 server/internal/models/monitor.go diff --git a/server/internal/checker/checker.go b/server/internal/checker/checker.go new file mode 100644 index 0000000..c4d49dd --- /dev/null +++ b/server/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/server/internal/models/monitor.go b/server/internal/models/monitor.go new file mode 100644 index 0000000..78b9978 --- /dev/null +++ b/server/internal/models/monitor.go @@ -0,0 +1,77 @@ +package models + +import ( + "time" + + "go.mongodb.org/mongo-driver/v2/bson" +) + +// Monitor check types. +const ( + MonitorHTTP = "http" + MonitorTCP = "tcp" + MonitorICMP = "icmp" + MonitorTLS = "tls" +) + +// Monitor status values. +const ( + StatusUp = "up" + StatusDown = "down" + StatusPending = "pending" +) + +// RunnerServer is the reserved Runner value for server-run monitors. Any other +// value is treated as a server_id whose agent runs the check locally. +const RunnerServer = "server" + +type MonitorTarget struct { + URL string `bson:"url,omitempty" json:"url,omitempty"` + Host string `bson:"host,omitempty" json:"host,omitempty"` + Port int `bson:"port,omitempty" json:"port,omitempty"` + Method string `bson:"method,omitempty" json:"method,omitempty"` + ExpectedStatus int `bson:"expected_status,omitempty" json:"expected_status,omitempty"` + Keyword string `bson:"keyword,omitempty" json:"keyword,omitempty"` + TLSWarnDays int `bson:"tls_warn_days,omitempty" json:"tls_warn_days,omitempty"` +} + +type MonitorState struct { + Status string `bson:"status" json:"status"` // up|down|pending + LastCheckAt *time.Time `bson:"last_check_at,omitempty" json:"last_check_at,omitempty"` + LatencyMs int `bson:"latency_ms" json:"latency_ms"` + Message string `bson:"message,omitempty" json:"message,omitempty"` + CertExpiryAt *time.Time `bson:"cert_expiry_at,omitempty" json:"cert_expiry_at,omitempty"` + Fails int `bson:"fails" json:"fails"` // consecutive failures + LastNotifiedAt *time.Time `bson:"last_notified_at,omitempty" json:"last_notified_at,omitempty"` +} + +type Monitor struct { + ID bson.ObjectID `bson:"_id,omitempty" json:"_id,omitempty"` + MonitorID string `bson:"monitor_id" json:"monitor_id"` + Name string `bson:"name" json:"name"` + Type string `bson:"type" json:"type"` // http|tcp|icmp|tls + Target MonitorTarget `bson:"target" json:"target"` + IntervalSec int `bson:"interval_sec" json:"interval_sec"` + Runner string `bson:"runner" json:"runner"` // "server" or a server_id + Retries int `bson:"retries" json:"retries"` // consecutive fails before down + Enabled bool `bson:"enabled" json:"enabled"` + ChannelIDs []string `bson:"channel_ids,omitempty" json:"channel_ids,omitempty"` + State MonitorState `bson:"state" json:"state"` + CreatedAt time.Time `bson:"created_at" json:"created_at"` +} + +type Incident struct { + IncidentID string `bson:"incident_id" json:"incident_id"` + MonitorID string `bson:"monitor_id" json:"monitor_id"` + StartedAt time.Time `bson:"started_at" json:"started_at"` + ResolvedAt *time.Time `bson:"resolved_at,omitempty" json:"resolved_at,omitempty"` + Cause string `bson:"cause,omitempty" json:"cause,omitempty"` +} + +type Rollup struct { + MonitorID string `bson:"monitor_id" json:"monitor_id"` + PeriodStart time.Time `bson:"period_start" json:"period_start"` // hour bucket + Checks int `bson:"checks" json:"checks"` + UpCount int `bson:"up_count" json:"up_count"` + SumLatency int64 `bson:"sum_latency" json:"sum_latency"` +}