diff --git a/server/cmd/main.go b/server/cmd/main.go index 4aa877f..cec38b2 100644 --- a/server/cmd/main.go +++ b/server/cmd/main.go @@ -11,6 +11,7 @@ import ( "github.com/mrhid6/vantage/server/internal/auth" "github.com/mrhid6/vantage/server/internal/db" grpcserver "github.com/mrhid6/vantage/server/internal/grpc" + "github.com/mrhid6/vantage/server/internal/monitorsched" "github.com/mrhid6/vantage/server/internal/services" ) @@ -67,6 +68,9 @@ func main() { } }() + // Start the server-side monitor scheduler. + monitorsched.Start(context.Background()) + // Start REST server r := gin.New() r.Use(gin.Recovery()) diff --git a/server/internal/api/handlers.go b/server/internal/api/handlers.go index 52a0174..2b10a31 100644 --- a/server/internal/api/handlers.go +++ b/server/internal/api/handlers.go @@ -80,6 +80,7 @@ func RegisterRoutes(r *gin.Engine) { apiGroup.GET("/console/tunnel", consoleTunnel) registerWorkflowRoutes(apiGroup) + registerMonitorRoutes(apiGroup) } } diff --git a/server/internal/api/monitors.go b/server/internal/api/monitors.go new file mode 100644 index 0000000..2e5f614 --- /dev/null +++ b/server/internal/api/monitors.go @@ -0,0 +1,139 @@ +package api + +import ( + "net/http" + "time" + + "github.com/gin-gonic/gin" + "github.com/mrhid6/vantage/server/internal/models" + "github.com/mrhid6/vantage/server/internal/services" + "go.mongodb.org/mongo-driver/v2/bson" +) + +func registerMonitorRoutes(g *gin.RouterGroup) { + g.GET("/monitors", listMonitors) + g.POST("/monitors", createMonitor) + g.GET("/monitors/:id", getMonitor) + g.PUT("/monitors/:id", updateMonitor) + g.DELETE("/monitors/:id", deleteMonitor) + g.GET("/monitors/:id/incidents", getMonitorIncidents) + g.GET("/monitors/:id/uptime", getMonitorUptime) +} + +func listMonitors(c *gin.Context) { + monitors, err := services.ListMonitors() + if err != nil { + c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) + return + } + c.JSON(http.StatusOK, monitors) +} + +func createMonitor(c *gin.Context) { + var m models.Monitor + if err := c.ShouldBindJSON(&m); err != nil { + c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()}) + return + } + if m.Name == "" || m.Type == "" { + c.JSON(http.StatusBadRequest, gin.H{"error": "name and type are required"}) + return + } + created, err := services.CreateMonitor(&m) + if err != nil { + c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) + return + } + c.JSON(http.StatusCreated, created) +} + +func getMonitor(c *gin.Context) { + m, err := services.GetMonitor(c.Param("id")) + if err != nil { + c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) + return + } + if m == nil { + c.JSON(http.StatusNotFound, gin.H{"error": "monitor not found"}) + return + } + c.JSON(http.StatusOK, m) +} + +func updateMonitor(c *gin.Context) { + var body struct { + Name *string `json:"name"` + Type *string `json:"type"` + Target *models.MonitorTarget `json:"target"` + IntervalSec *int `json:"interval_sec"` + Runner *string `json:"runner"` + Retries *int `json:"retries"` + Enabled *bool `json:"enabled"` + ChannelIDs *[]string `json:"channel_ids"` + } + if err := c.ShouldBindJSON(&body); err != nil { + c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()}) + return + } + upd := bson.M{} + if body.Name != nil { + upd["name"] = *body.Name + } + if body.Type != nil { + upd["type"] = *body.Type + } + if body.Target != nil { + upd["target"] = *body.Target + } + if body.IntervalSec != nil { + upd["interval_sec"] = *body.IntervalSec + } + if body.Runner != nil { + upd["runner"] = *body.Runner + } + if body.Retries != nil { + upd["retries"] = *body.Retries + } + if body.Enabled != nil { + upd["enabled"] = *body.Enabled + } + if body.ChannelIDs != nil { + upd["channel_ids"] = *body.ChannelIDs + } + if len(upd) == 0 { + c.JSON(http.StatusBadRequest, gin.H{"error": "no fields to update"}) + return + } + if err := services.UpdateMonitor(c.Param("id"), upd); err != nil { + c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) + return + } + c.Status(http.StatusNoContent) +} + +func deleteMonitor(c *gin.Context) { + if err := services.DeleteMonitor(c.Param("id")); err != nil { + c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) + return + } + c.Status(http.StatusNoContent) +} + +func getMonitorIncidents(c *gin.Context) { + incidents, err := services.ListIncidents(c.Param("id"), 50) + if err != nil { + c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) + return + } + c.JSON(http.StatusOK, incidents) +} + +func getMonitorUptime(c *gin.Context) { + since := time.Now().Add(-30 * 24 * time.Hour) + rollups, err := services.UptimeRollups(c.Param("id"), since) + if err != nil { + c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) + return + } + c.JSON(http.StatusOK, rollups) +} diff --git a/server/internal/monitorsched/scheduler.go b/server/internal/monitorsched/scheduler.go new file mode 100644 index 0000000..f3b33a8 --- /dev/null +++ b/server/internal/monitorsched/scheduler.go @@ -0,0 +1,107 @@ +// Package monitorsched runs server-side monitors on their configured interval +// and funnels results through services.IngestResult. Agent-run monitors +// (runner != "server") are excluded — those execute on the agent. +package monitorsched + +import ( + "context" + "log" + "sync" + "time" + + "github.com/mrhid6/vantage/server/internal/checker" + "github.com/mrhid6/vantage/server/internal/models" + "github.com/mrhid6/vantage/server/internal/services" +) + +// reloadInterval controls how often the scheduler re-reads monitor definitions +// so CRUD changes (new/removed/edited monitors) take effect. +const reloadInterval = 30 * time.Second + +type runner struct { + monitorID string + intervalSec int + cancel context.CancelFunc +} + +// Start launches the scheduler loop. It returns immediately; the loop runs until +// ctx is cancelled. +func Start(ctx context.Context) { + go loop(ctx) +} + +func loop(ctx context.Context) { + active := map[string]*runner{} + var mu sync.Mutex + + sync := func() { + monitors, err := services.ListMonitorsForRunner(models.RunnerServer) + if err != nil { + log.Printf("monitorsched: list monitors: %v", err) + return + } + want := map[string]models.Monitor{} + for _, m := range monitors { + want[m.MonitorID] = m + } + + mu.Lock() + defer mu.Unlock() + // Stop runners for monitors that vanished or changed interval. + for id, r := range active { + m, ok := want[id] + if !ok || m.IntervalSec != r.intervalSec { + r.cancel() + delete(active, id) + } + } + // Start runners for new/changed monitors. + for id, m := range want { + if _, ok := active[id]; ok { + continue + } + rctx, cancel := context.WithCancel(ctx) + active[id] = &runner{monitorID: id, intervalSec: m.IntervalSec, cancel: cancel} + go runMonitor(rctx, m) + } + } + + sync() + t := time.NewTicker(reloadInterval) + defer t.Stop() + for { + select { + case <-ctx.Done(): + return + case <-t.C: + sync() + } + } +} + +func runMonitor(ctx context.Context, m models.Monitor) { + interval := time.Duration(m.IntervalSec) * time.Second + if interval <= 0 { + interval = 60 * time.Second + } + spec := services.SpecFor(&m) + + run := func() { + res := checker.Run(ctx, spec) + if err := services.IngestResult(m.MonitorID, res); err != nil { + log.Printf("monitorsched: ingest %s: %v", m.MonitorID, err) + } + } + + run() // check immediately on (re)start + t := time.NewTicker(interval) + defer t.Stop() + for { + select { + case <-ctx.Done(): + return + case <-t.C: + run() + } + } +}