From 8f3a27100f843ddbe8c9cae4da2d2a5fc22c3044 Mon Sep 17 00:00:00 2001 From: mrhid6 Date: Tue, 21 Jul 2026 14:22:55 +0100 Subject: [PATCH] feat(server): multi-channel monitor notifications --- server/internal/api/channels.go | 96 +++++++++++++++++++++++++ server/internal/api/handlers.go | 1 + server/internal/models/channel.go | 29 ++++++++ server/internal/notify/dispatch.go | 64 +++++++++++++++++ server/internal/notify/http.go | 71 +++++++++++++++++++ server/internal/notify/smtp.go | 45 ++++++++++++ server/internal/services/channels.go | 100 +++++++++++++++++++++++++++ server/internal/services/monitors.go | 36 ++++++++-- 8 files changed, 438 insertions(+), 4 deletions(-) create mode 100644 server/internal/api/channels.go create mode 100644 server/internal/models/channel.go create mode 100644 server/internal/notify/dispatch.go create mode 100644 server/internal/notify/http.go create mode 100644 server/internal/notify/smtp.go create mode 100644 server/internal/services/channels.go diff --git a/server/internal/api/channels.go b/server/internal/api/channels.go new file mode 100644 index 0000000..f707904 --- /dev/null +++ b/server/internal/api/channels.go @@ -0,0 +1,96 @@ +package api + +import ( + "net/http" + + "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 registerChannelRoutes(g *gin.RouterGroup) { + g.GET("/channels", listChannels) + g.POST("/channels", createChannel) + g.PUT("/channels/:id", updateChannel) + g.DELETE("/channels/:id", deleteChannel) + g.POST("/channels/:id/test", testChannel) +} + +func listChannels(c *gin.Context) { + channels, err := services.ListChannels() + if err != nil { + c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) + return + } + c.JSON(http.StatusOK, channels) +} + +func createChannel(c *gin.Context) { + var ch models.NotificationChannel + if err := c.ShouldBindJSON(&ch); err != nil { + c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()}) + return + } + if ch.Name == "" || ch.Type == "" { + c.JSON(http.StatusBadRequest, gin.H{"error": "name and type are required"}) + return + } + created, err := services.CreateChannel(&ch) + if err != nil { + c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) + return + } + c.JSON(http.StatusCreated, created) +} + +func updateChannel(c *gin.Context) { + var body struct { + Name *string `json:"name"` + Type *string `json:"type"` + Config *map[string]string `json:"config"` + Enabled *bool `json:"enabled"` + } + 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.Config != nil { + upd["config"] = *body.Config + } + if body.Enabled != nil { + upd["enabled"] = *body.Enabled + } + if len(upd) == 0 { + c.JSON(http.StatusBadRequest, gin.H{"error": "no fields to update"}) + return + } + if err := services.UpdateChannel(c.Param("id"), upd); err != nil { + c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) + return + } + c.Status(http.StatusNoContent) +} + +func deleteChannel(c *gin.Context) { + if err := services.DeleteChannel(c.Param("id")); err != nil { + c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) + return + } + c.Status(http.StatusNoContent) +} + +func testChannel(c *gin.Context) { + if err := services.TestChannel(c.Param("id")); err != nil { + c.JSON(http.StatusBadGateway, gin.H{"error": err.Error()}) + return + } + c.JSON(http.StatusOK, gin.H{"status": "sent"}) +} diff --git a/server/internal/api/handlers.go b/server/internal/api/handlers.go index 2b10a31..7132c12 100644 --- a/server/internal/api/handlers.go +++ b/server/internal/api/handlers.go @@ -81,6 +81,7 @@ func RegisterRoutes(r *gin.Engine) { registerWorkflowRoutes(apiGroup) registerMonitorRoutes(apiGroup) + registerChannelRoutes(apiGroup) } } diff --git a/server/internal/models/channel.go b/server/internal/models/channel.go new file mode 100644 index 0000000..c166cef --- /dev/null +++ b/server/internal/models/channel.go @@ -0,0 +1,29 @@ +package models + +import ( + "time" + + "go.mongodb.org/mongo-driver/v2/bson" +) + +// Notification channel types. +const ( + ChannelWebhook = "webhook" + ChannelSMTP = "smtp" + ChannelDiscord = "discord" + ChannelSlack = "slack" + ChannelTelegram = "telegram" +) + +// NotificationChannel is an outbound alert destination. Config holds +// type-specific settings (e.g. url; or smtp host/port/username/password/from/to; +// or telegram token/chat_id). +type NotificationChannel struct { + ID bson.ObjectID `bson:"_id,omitempty" json:"_id,omitempty"` + ChannelID string `bson:"channel_id" json:"channel_id"` + Name string `bson:"name" json:"name"` + Type string `bson:"type" json:"type"` + Config map[string]string `bson:"config" json:"config"` + Enabled bool `bson:"enabled" json:"enabled"` + CreatedAt time.Time `bson:"created_at" json:"created_at"` +} diff --git a/server/internal/notify/dispatch.go b/server/internal/notify/dispatch.go new file mode 100644 index 0000000..00a846b --- /dev/null +++ b/server/internal/notify/dispatch.go @@ -0,0 +1,64 @@ +// Package notify formats and delivers monitor state-change alerts to +// notification channels. It depends only on models so services can call it +// without an import cycle. +package notify + +import ( + "fmt" + "time" + + "github.com/mrhid6/vantage/server/internal/models" +) + +// Event describes a monitor state transition worth alerting on. +type Event struct { + MonitorName string + Type string + OldStatus string + NewStatus string + Message string + Time time.Time +} + +// title is a short one-line summary used by the text-based channels. +func (e Event) title() string { + verb := "recovered" + if e.NewStatus == models.StatusDown { + verb = "is DOWN" + } + s := fmt.Sprintf("[Vantage] %s (%s) %s", e.MonitorName, e.Type, verb) + if e.Message != "" { + s += ": " + e.Message + } + return s +} + +// Dispatch delivers ev to a single channel, formatting per channel type. +func Dispatch(ch models.NotificationChannel, ev Event) error { + switch ch.Type { + case models.ChannelWebhook: + return dispatchWebhook(ch, ev) + case models.ChannelDiscord: + return dispatchDiscord(ch, ev) + case models.ChannelSlack: + return dispatchSlack(ch, ev) + case models.ChannelTelegram: + return dispatchTelegram(ch, ev) + case models.ChannelSMTP: + return dispatchSMTP(ch, ev) + default: + return fmt.Errorf("unknown channel type: %s", ch.Type) + } +} + +// Test delivers a synthetic event so users can verify a channel's configuration. +func Test(ch models.NotificationChannel) error { + return Dispatch(ch, Event{ + MonitorName: "Test monitor", + Type: "http", + OldStatus: models.StatusUp, + NewStatus: models.StatusDown, + Message: "this is a test alert from Vantage", + Time: time.Now(), + }) +} diff --git a/server/internal/notify/http.go b/server/internal/notify/http.go new file mode 100644 index 0000000..5880fec --- /dev/null +++ b/server/internal/notify/http.go @@ -0,0 +1,71 @@ +package notify + +import ( + "bytes" + "encoding/json" + "fmt" + "net/http" + "time" + + "github.com/mrhid6/vantage/server/internal/models" +) + +var httpClient = &http.Client{Timeout: 10 * time.Second} + +func postJSON(target string, payload any) error { + body, err := json.Marshal(payload) + if err != nil { + return err + } + resp, err := httpClient.Post(target, "application/json", bytes.NewReader(body)) + if err != nil { + return err + } + defer resp.Body.Close() + if resp.StatusCode >= 300 { + return fmt.Errorf("HTTP %d from %s", resp.StatusCode, target) + } + return nil +} + +// dispatchWebhook posts the full event as JSON to a user-supplied URL. +func dispatchWebhook(ch models.NotificationChannel, ev Event) error { + target := ch.Config["url"] + if target == "" { + return fmt.Errorf("webhook: missing url") + } + return postJSON(target, map[string]any{ + "monitor": ev.MonitorName, + "type": ev.Type, + "old_status": ev.OldStatus, + "new_status": ev.NewStatus, + "message": ev.Message, + "time": ev.Time.Format(time.RFC3339), + }) +} + +func dispatchDiscord(ch models.NotificationChannel, ev Event) error { + target := ch.Config["url"] + if target == "" { + return fmt.Errorf("discord: missing url") + } + return postJSON(target, map[string]string{"content": ev.title()}) +} + +func dispatchSlack(ch models.NotificationChannel, ev Event) error { + target := ch.Config["url"] + if target == "" { + return fmt.Errorf("slack: missing url") + } + return postJSON(target, map[string]string{"text": ev.title()}) +} + +func dispatchTelegram(ch models.NotificationChannel, ev Event) error { + token := ch.Config["token"] + chatID := ch.Config["chat_id"] + if token == "" || chatID == "" { + return fmt.Errorf("telegram: missing token or chat_id") + } + api := fmt.Sprintf("https://api.telegram.org/bot%s/sendMessage", token) + return postJSON(api, map[string]string{"chat_id": chatID, "text": ev.title()}) +} diff --git a/server/internal/notify/smtp.go b/server/internal/notify/smtp.go new file mode 100644 index 0000000..481799b --- /dev/null +++ b/server/internal/notify/smtp.go @@ -0,0 +1,45 @@ +package notify + +import ( + "fmt" + "net/smtp" + "strings" + + "github.com/mrhid6/vantage/server/internal/models" +) + +// dispatchSMTP sends the alert as a plain-text email. Config keys: host, port, +// username, password, from, to. Auth is skipped when username is empty. +func dispatchSMTP(ch models.NotificationChannel, ev Event) error { + host := ch.Config["host"] + port := ch.Config["port"] + from := ch.Config["from"] + to := ch.Config["to"] + if host == "" || port == "" || from == "" || to == "" { + return fmt.Errorf("smtp: missing host/port/from/to") + } + + title := ev.title() + msg := strings.Join([]string{ + "From: " + from, + "To: " + to, + "Subject: " + title, + "", + title, + "", + "Monitor: " + ev.MonitorName, + "Status: " + ev.OldStatus + " -> " + ev.NewStatus, + "Time: " + ev.Time.String(), + }, "\r\n") + + var auth smtp.Auth + if user := ch.Config["username"]; user != "" { + auth = smtp.PlainAuth("", user, ch.Config["password"], host) + } + addr := host + ":" + port + recipients := strings.Split(to, ",") + for i := range recipients { + recipients[i] = strings.TrimSpace(recipients[i]) + } + return smtp.SendMail(addr, auth, from, recipients, []byte(msg)) +} diff --git a/server/internal/services/channels.go b/server/internal/services/channels.go new file mode 100644 index 0000000..d37a93b --- /dev/null +++ b/server/internal/services/channels.go @@ -0,0 +1,100 @@ +package services + +import ( + "errors" + "time" + + "github.com/google/uuid" + "github.com/mrhid6/vantage/server/internal/db" + "github.com/mrhid6/vantage/server/internal/models" + "github.com/mrhid6/vantage/server/internal/notify" + "go.mongodb.org/mongo-driver/v2/bson" + "go.mongodb.org/mongo-driver/v2/mongo" + "go.mongodb.org/mongo-driver/v2/mongo/options" +) + +func ListChannels() ([]models.NotificationChannel, error) { + ctx, cancel := monCtx() + defer cancel() + cur, err := db.Col("notification_channels").Find(ctx, bson.M{}, options.Find().SetSort(bson.M{"created_at": 1})) + if err != nil { + return nil, err + } + var out []models.NotificationChannel + if err := cur.All(ctx, &out); err != nil { + return nil, err + } + return out, nil +} + +func GetChannel(channelID string) (*models.NotificationChannel, error) { + ctx, cancel := monCtx() + defer cancel() + var ch models.NotificationChannel + err := db.Col("notification_channels").FindOne(ctx, bson.M{"channel_id": channelID}).Decode(&ch) + if errors.Is(err, mongo.ErrNoDocuments) { + return nil, nil + } + if err != nil { + return nil, err + } + return &ch, nil +} + +// GetChannels loads multiple channels by ID, skipping any not found. +func GetChannels(channelIDs []string) ([]models.NotificationChannel, error) { + if len(channelIDs) == 0 { + return nil, nil + } + ctx, cancel := monCtx() + defer cancel() + cur, err := db.Col("notification_channels").Find(ctx, bson.M{"channel_id": bson.M{"$in": channelIDs}}) + if err != nil { + return nil, err + } + var out []models.NotificationChannel + if err := cur.All(ctx, &out); err != nil { + return nil, err + } + return out, nil +} + +func CreateChannel(ch *models.NotificationChannel) (*models.NotificationChannel, error) { + ctx, cancel := monCtx() + defer cancel() + ch.ChannelID = uuid.NewString() + ch.CreatedAt = time.Now() + if ch.Config == nil { + ch.Config = map[string]string{} + } + if _, err := db.Col("notification_channels").InsertOne(ctx, ch); err != nil { + return nil, err + } + return ch, nil +} + +func UpdateChannel(channelID string, upd bson.M) error { + ctx, cancel := monCtx() + defer cancel() + _, err := db.Col("notification_channels").UpdateOne(ctx, bson.M{"channel_id": channelID}, bson.M{"$set": upd}) + return err +} + +func DeleteChannel(channelID string) error { + ctx, cancel := monCtx() + defer cancel() + _, err := db.Col("notification_channels").DeleteOne(ctx, bson.M{"channel_id": channelID}) + return err +} + +// TestChannel sends a synthetic alert to verify configuration. +func TestChannel(channelID string) error { + ch, err := GetChannel(channelID) + if err != nil { + return err + } + if ch == nil { + return errors.New("channel not found") + } + return notify.Test(*ch) +} diff --git a/server/internal/services/monitors.go b/server/internal/services/monitors.go index da00683..992d414 100644 --- a/server/internal/services/monitors.go +++ b/server/internal/services/monitors.go @@ -3,12 +3,14 @@ package services import ( "context" "errors" + "log" "time" "github.com/google/uuid" "github.com/mrhid6/vantage/server/internal/checker" "github.com/mrhid6/vantage/server/internal/db" "github.com/mrhid6/vantage/server/internal/models" + "github.com/mrhid6/vantage/server/internal/notify" "go.mongodb.org/mongo-driver/v2/bson" "go.mongodb.org/mongo-driver/v2/mongo" "go.mongodb.org/mongo-driver/v2/mongo/options" @@ -233,9 +235,35 @@ func IngestResult(monitorID string, res checker.Result) error { return nil } -// notifyTransition dispatches notifications on an up<->down transition. Wired up -// in Task 13 (P3); a no-op until then. +// notifyTransition dispatches notifications on an up<->down transition to each +// enabled channel bound to the monitor. Deliveries run in the background; +// failures are logged, not fatal. func notifyTransition(m *models.Monitor, newStatus, message string) { - // TODO(P3): resolve m.ChannelIDs and dispatch via the notify package, - // honouring a resend interval tracked on state.last_notified_at. + if len(m.ChannelIDs) == 0 { + return + } + channels, err := GetChannels(m.ChannelIDs) + if err != nil { + log.Printf("notify: load channels for %s: %v", m.MonitorID, err) + return + } + ev := notify.Event{ + MonitorName: m.Name, + Type: m.Type, + OldStatus: m.State.Status, + NewStatus: newStatus, + Message: message, + Time: time.Now(), + } + for _, ch := range channels { + if !ch.Enabled { + continue + } + go func(c models.NotificationChannel) { + if err := notify.Dispatch(c, ev); err != nil { + log.Printf("notify: dispatch to %s (%s): %v", c.Name, c.Type, err) + } + }(ch) + } + _ = UpdateMonitor(m.MonitorID, bson.M{"state.last_notified_at": time.Now()}) }