Files
2026-09-16 06:26:16 +02:00

196 lines
5.7 KiB
Go

package gateway
import (
"context"
"encoding/json"
"errors"
"fmt"
"regexp"
"strings"
"text/template"
"unicode/utf16"
"github.com/example/notify-gateway/internal/config"
"github.com/example/notify-gateway/internal/divera"
"github.com/example/notify-gateway/internal/model"
"github.com/example/notify-gateway/internal/outbound"
)
type Dispatcher struct {
store *config.Store
divera *divera.Client
}
func New(store *config.Store, client *divera.Client) *Dispatcher {
return &Dispatcher{store: store, divera: client}
}
func (d *Dispatcher) Dispatch(ctx context.Context, msg model.InboundMessage) ([]model.DeliveryResult, error) {
cfg := d.store.Get()
var results []model.DeliveryResult
matched := false
var failures []error
for _, m := range cfg.Mappings {
if !m.Enabled || !matches(m, msg) {
continue
}
matched = true
r := model.DeliveryResult{MappingID: m.ID, MappingName: m.Name, Target: m.Target, OutboundID: m.OutboundID}
payload, kind, err := buildPayload(m, msg)
if err == nil {
if config.IsOutbound(kind) {
var destination *config.OutboundConfig
for _, o := range cfg.Outbounds {
if o.ID == m.OutboundID && o.Provider == kind {
destination = &o
break
}
}
if destination == nil {
err = fmt.Errorf("outbound destination missing or provider mismatch")
} else {
var resp outbound.Response
resp, err = outbound.Send(ctx, *destination, payload)
r.StatusCode, r.Response = resp.StatusCode, resp.Body
}
} else {
var resp divera.Response
resp, err = d.divera.Create(ctx, kind, payload)
r.StatusCode, r.Response = resp.StatusCode, string(resp.Body)
}
}
if err != nil {
r.Error = err.Error()
failures = append(failures, fmt.Errorf("mapping %s: %w", m.ID, err))
}
results = append(results, r)
}
if !matched {
return results, fmt.Errorf("no mapping matched source=%s channel=%s", msg.Source, msg.Channel)
}
return results, errors.Join(failures...)
}
func matches(m config.Mapping, msg model.InboundMessage) bool {
if m.Source != "" && m.Source != "any" && m.Source != msg.Source {
return false
}
if m.MinPriority != 0 && msg.Priority < m.MinPriority {
return false
}
for _, condition := range [][2]string{{m.ChannelRegex, msg.Channel}, {m.TitleRegex, msg.Title}, {m.MessageRegex, msg.Message}} {
pat, value := condition[0], condition[1]
if pat == "" {
continue
}
ok, err := regexp.MatchString(pat, value)
if err != nil || !ok {
return false
}
}
return true
}
func render(s string, msg model.InboundMessage) (string, error) {
if s == "" {
return "", nil
}
t, err := template.New("m").Option("missingkey=zero").Parse(s)
if err != nil {
return "", err
}
var b strings.Builder
if err := t.Execute(&b, msg); err != nil {
return "", err
}
return b.String(), nil
}
func buildPayload(m config.Mapping, msg model.InboundMessage) (map[string]any, string, error) {
title, err := render(m.TitleTemplate, msg)
if err != nil {
return nil, "", err
}
text, err := render(m.TextTemplate, msg)
if err != nil {
return nil, "", err
}
address, err := render(m.AddressTemplate, msg)
if err != nil {
return nil, "", err
}
if address == "" {
address = msg.Address
}
switch strings.ToLower(m.Target) {
case "smtp", "ntfy", "gotify":
if strings.TrimSpace(text) == "" {
return nil, "", fmt.Errorf("message must not be empty")
}
return map[string]any{"title": title, "message": text, "priority": msg.Priority}, strings.ToLower(m.Target), nil
case "discord":
content := strings.TrimSpace(title + "\n" + text)
if content == "" {
return nil, "", fmt.Errorf("Discord content must not be empty")
}
if len(utf16.Encode([]rune(content))) > 2000 {
return nil, "", fmt.Errorf("Discord content exceeds 2000 characters; shorten title/text templates")
}
return map[string]any{"content": content, "allowed_mentions": map[string]any{"parse": []string{}}}, "discord", nil
case "webhook":
return map[string]any{"source": msg.Source, "channel": msg.Channel, "title": title, "message": text, "address": address, "priority": msg.Priority, "tags": msg.Tags, "received_at": msg.ReceivedAt}, "webhook", nil
}
obj := map[string]any{"title": title, "text": text, "notification_type": m.NotificationType, "send_push": m.SendPush, "send_sms": m.SendSMS, "send_call": m.SendCall, "send_mail": m.SendMail, "send_pager": m.SendPager, "private_mode": m.PrivateMode}
if address != "" {
obj["address"] = address
}
if len(m.ClusterRoutes) > 0 {
routes := make(map[string]any, len(m.ClusterRoutes))
for id, notificationType := range m.ClusterRoutes {
routes[id] = map[string]any{"notification_type": notificationType}
}
obj["cluster"] = routes
} else if len(m.Clusters) > 0 {
obj["cluster"] = m.Clusters
}
if len(m.Groups) > 0 {
obj["group"] = m.Groups
}
if len(m.Users) > 0 {
obj["user_cluster_relation"] = m.Users
}
if len(m.Vehicles) > 0 {
obj["vehicle"] = m.Vehicles
}
for k, v := range m.Extra {
obj[k] = v
}
var root, kind string
switch strings.ToLower(m.Target) {
case "alarm", "alarms":
root, kind = "Alarm", "alarms"
case "news", "message", "mitteilung":
root, kind = "News", "news"
case "event", "termin":
root, kind = "Event", "events"
default:
return nil, "", fmt.Errorf("unsupported target %q", m.Target)
}
return map[string]any{root: obj}, kind, nil
}
func Pretty(v any) string { b, _ := json.MarshalIndent(v, "", " "); return string(b) }
// Preview evaluates a single mapping against a sample message without sending anything.
func Preview(m config.Mapping, msg model.InboundMessage) (map[string]any, string, bool, error) {
matched := matches(m, msg)
if !matched {
return nil, "", false, nil
}
payload, kind, err := buildPayload(m, msg)
if err != nil {
return nil, "", true, err
}
return payload, kind, true, nil
}