Files
mmnotify/main.go
T
elton 509475e17a mmnotify: der Rueckweg aus dem Teamspace, als Werkzeug
Andreas: "kannst du neben dem goimapnotify einen mmnotify daneben setzen?
das waere doch genau das werkzeug, was wir nun fuer alle brauchen. Auch
gerne in go, auch gerne sehr robust, automatischer wiederaufbau und so."

Dieselbe Bauform wie goimapnotify: ein kleiner Prozess, der lauscht und
bei einem Ereignis einen Befehl ausfuehrt. Er schreibt nie selbst in
einen Kanal -- er klopft nur an, der Agent liest und antwortet mit seinen
eigenen Werkzeugen.

Zwei Wege, und der Boden ist der tragende. Der Websocket beschleunigt,
darunter laeuft ein Boden, der ohnehin nachsieht. Uebernommen aus hive
(channels/team.py) und dort begruendet: ein Strom, der still ausfaellt,
darf uns spaet machen, nie blind.

Hier ist das nicht nur Vorsicht: Gegen Mattermost 11.7.10 gemessen kommt
die Verbindung zustande, die Anmeldung wird mit status OK quittiert, und
danach kommen hello und status_change -- aber KEIN posted, auch wenn ein
anderes Konto in einen Kanal schreibt, in dem der Agent nachweislich
Mitglied ist. Ursache offen, Verdacht sind Broadcasts an Verbindungen,
die ueber ein Personal Access Token angemeldet sind. Mit dem Boden
arbeitet der Dienst trotzdem, nur langsamer.

Daraus folgt die zweite Entscheidung, ebenfalls von hive: mmnotify liest
den Beitrag NICHT aus dem Ereignis. Zwei Leser waeren zwei Antworten.

Robustheit, wo sie noetig ist: Wiederaufbau mit wachsendem Abstand bis
60 s; Herzschlag alle 30 s gegen halboffene Verbindungen; hoechstens ein
Lauf gleichzeitig, sonst loesen fuenf Beitraege fuenf Agentenlaeufe aus,
die sich ueberholen; eigene Beitraege und system-Ereignisse wecken nie.
2026-09-01 20:38:01 +02:00

348 lines
10 KiB
Go

// mmnotify — hört am Mattermost-Websocket und weckt einen Agenten, wenn er
// gemeint ist.
//
// Gegenstück zu goimapnotify für Mail: ein kleiner Prozess, der lauscht und
// bei einem Ereignis einen Befehl ausführt. Der Agent selbst liest und
// antwortet danach über seine eigenen Werkzeuge — mmnotify schreibt nie in
// einen Kanal, es klopft nur an.
//
// Aufruf: mmnotify -conf /etc/mmnotify.json
//
// ZWEI WEGE, UND DER BODEN IST DER TRAGENDE.
//
// Der Websocket ist der Beschleuniger: Kommt ein Ereignis, wird sofort geweckt.
// Darunter liegt ein Boden (floor_seconds), der ohnehin läuft. Diese Aufteilung
// ist von hive übernommen und dort begründet: "A stream that dies quietly must
// make us LATE, never blind." Ein Strom, der still ausfällt, kostet dann
// Latenz statt Nachrichten.
//
// Das ist hier keine Vorsichtsmaßnahme, sondern nötig. Gemessen am 2026-09-01
// gegen team.42i.org (Mattermost 11.7.10): Die Verbindung kommt zustande, die
// Anmeldung wird mit "status: OK" quittiert, und danach kommen `hello` und
// `status_change` — aber KEIN `posted`, auch wenn ein anderes Konto in einen
// Kanal schreibt, in dem der Agent nachweislich Mitglied ist. Ursache noch
// offen (Verdacht: Broadcasts an Verbindungen, die über ein Personal Access
// Token angemeldet sind). Mit dem Boden funktioniert der Dienst trotzdem.
//
// Deshalb gilt auch: mmnotify liest den Beitrag NICHT aus dem Ereignis. Es
// klopft nur an; was zu tun ist, entscheidet der Agent, wenn er nachsieht.
// Zwei Leser wären zwei Antworten (hive, channels/team.py).
package main
import (
"encoding/json"
"flag"
"fmt"
"log"
"net/http"
"net/url"
"os"
"os/exec"
"os/signal"
"strings"
"sync"
"syscall"
"time"
"github.com/gorilla/websocket"
)
type Config struct {
// Basisadresse des Servers, z. B. https://team.42i.org
URL string `json:"url"`
// Persönliches Zugriffstoken des Agenten. Steht in der Konfiguration, weil
// das Werkzeug sonst ein zweites Geheimnissystem bräuchte; die Datei gehört
// mit 0600 dem Dienstbenutzer.
Token string `json:"token"`
// Was ausgeführt wird, wenn der Agent gemeint ist. Bekommt die Ereignisdaten
// zusätzlich über die Umgebung (MM_CHANNEL, MM_CHANNEL_NAME, MM_SENDER,
// MM_POST_ID, MM_ROOT_ID, MM_MESSAGE).
OnMention string `json:"on_mention"`
// Optional: eigener Befehl für Direktnachrichten. Leer = OnMention.
OnDirect string `json:"on_direct"`
// Wie lange nach dem letzten Ereignis mindestens gewartet wird, bevor
// derselbe Befehl erneut läuft. Verhindert, dass fünf Beiträge in einer
// Minute fünf Agentenläufe auslösen.
DebounceSeconds int `json:"debounce_seconds"`
// Der Boden unter dem Stream: Auch ohne Ereignis wird nach so vielen
// Sekunden nachgesehen. Standard 60. Auf 0 zu setzen heißt, sich allein auf
// den Stream zu verlassen — siehe die Messung im Kopf dieser Datei.
FloorSeconds int `json:"floor_seconds"`
// Antwortet der Agent auf eigene Beiträge, entsteht eine Schleife. Der
// eigene Benutzer wird beim Start ermittelt und hier nicht gebraucht.
IgnoreUsers []string `json:"ignore_users"`
}
type Ereignis struct {
Event string `json:"event"`
Data map[string]interface{} `json:"data"`
Seq int64 `json:"seq"`
}
type Melder struct {
cfg Config
ich string // eigene user_id
ichName string
namen map[string]string
namenMu sync.Mutex
letzter time.Time
laufMu sync.Mutex
gesehen map[string]bool
}
func main() {
conf := flag.String("conf", "/etc/mmnotify.json", "Konfigurationsdatei")
flag.Parse()
roh, err := os.ReadFile(*conf)
if err != nil {
log.Fatalf("Konfiguration nicht lesbar: %v", err)
}
var cfg Config
if err := json.Unmarshal(roh, &cfg); err != nil {
log.Fatalf("Konfiguration ist kein gültiges JSON: %v", err)
}
cfg.URL = strings.TrimRight(cfg.URL, "/")
if cfg.URL == "" || cfg.Token == "" || cfg.OnMention == "" {
log.Fatal("url, token und on_mention sind Pflicht")
}
if cfg.DebounceSeconds == 0 {
cfg.DebounceSeconds = 5
}
if cfg.FloorSeconds == 0 {
cfg.FloorSeconds = 60
}
m := &Melder{cfg: cfg, namen: map[string]string{}, gesehen: map[string]bool{}}
if err := m.ichErmitteln(); err != nil {
log.Fatalf("Anmeldung fehlgeschlagen: %v", err)
}
log.Printf("mmnotify läuft als %s an %s", m.ichName, cfg.URL)
// Sauber beenden: Ein Dienst, der auf SIGTERM nicht reagiert, wird nach
// dem Timeout hart getötet und hinterlässt eine halboffene Verbindung.
schluss := make(chan os.Signal, 1)
signal.Notify(schluss, syscall.SIGINT, syscall.SIGTERM)
go func() {
<-schluss
log.Print("Beende auf Signal.")
os.Exit(0)
}()
if cfg.FloorSeconds > 0 {
go m.boden()
}
m.verbindungsschleife()
}
// verbindungsschleife hält die Verbindung offen und baut sie wieder auf.
// Der Backoff wächst bis eine Minute: Ein Server, der gerade neu startet, darf
// nicht von einem Client bestürmt werden, der im Sekundentakt anklopft.
func (m *Melder) verbindungsschleife() {
warte := time.Second
for {
err := m.lauschen()
if err != nil {
log.Printf("Verbindung verloren: %v — neuer Versuch in %s", err, warte)
}
time.Sleep(warte)
warte *= 2
if warte > time.Minute {
warte = time.Minute
}
// Nach einer erfolgreichen Verbindung zählt der Backoff wieder von
// vorn; das erledigt lauschen() über zuruecksetzen unten.
}
}
func (m *Melder) wsAdresse() string {
u, _ := url.Parse(m.cfg.URL)
if u.Scheme == "https" {
u.Scheme = "wss"
} else {
u.Scheme = "ws"
}
u.Path = "/api/v4/websocket"
return u.String()
}
func (m *Melder) lauschen() error {
c, _, err := websocket.DefaultDialer.Dial(m.wsAdresse(), http.Header{})
if err != nil {
return fmt.Errorf("Verbindungsaufbau: %w", err)
}
defer c.Close()
// Mattermost erwartet die Anmeldung als erste Nachricht auf dem Kanal.
if err := c.WriteJSON(map[string]interface{}{
"seq": 1, "action": "authentication_challenge",
"data": map[string]string{"token": m.cfg.Token},
}); err != nil {
return fmt.Errorf("Anmeldung senden: %w", err)
}
// Herzschlag: Eine halboffene Verbindung liest sich wie Stille. Der Ping
// zwingt sie, sich zu erklären.
c.SetReadDeadline(time.Now().Add(90 * time.Second))
c.SetPongHandler(func(string) error {
c.SetReadDeadline(time.Now().Add(90 * time.Second))
return nil
})
fertig := make(chan struct{})
defer close(fertig)
go func() {
t := time.NewTicker(30 * time.Second)
defer t.Stop()
for {
select {
case <-t.C:
if err := c.WriteControl(websocket.PingMessage, nil,
time.Now().Add(10*time.Second)); err != nil {
return
}
case <-fertig:
return
}
}
}()
log.Print("Verbunden, warte auf Ereignisse.")
for {
var ev Ereignis
if err := c.ReadJSON(&ev); err != nil {
return fmt.Errorf("Lesen: %w", err)
}
if ev.Event == "posted" {
m.beitrag(ev)
}
}
}
func (m *Melder) beitrag(ev Ereignis) {
rohPost, _ := ev.Data["post"].(string)
if rohPost == "" {
return
}
var p map[string]interface{}
if err := json.Unmarshal([]byte(rohPost), &p); err != nil {
return
}
postID, _ := p["id"].(string)
if postID == "" || m.gesehen[postID] {
return
}
m.gesehen[postID] = true
if len(m.gesehen) > 5000 {
m.gesehen = map[string]bool{postID: true}
}
autorID, _ := p["user_id"].(string)
if autorID == m.ich {
return // eigene Beiträge wecken niemanden
}
if typ, _ := p["type"].(string); strings.HasPrefix(typ, "system_") {
return // Beitritte, Umbenennungen, Kanal-Kopfzeilen
}
nachricht, _ := p["message"].(string)
kanalTyp, _ := ev.Data["channel_type"].(string)
kanalName, _ := ev.Data["channel_display_name"].(string)
kanalID, _ := p["channel_id"].(string)
autor, _ := ev.Data["sender_name"].(string)
autor = strings.TrimPrefix(autor, "@")
for _, u := range m.cfg.IgnoreUsers {
if strings.EqualFold(u, autor) {
return
}
}
direkt := kanalTyp == "D"
gemeint := direkt || strings.Contains(nachricht, "@"+m.ichName)
if !gemeint {
return // Mitlesen im Kanal ist kein Auftrag
}
wurzel, _ := p["root_id"].(string)
if wurzel == "" {
wurzel = postID
}
befehl := m.cfg.OnMention
if direkt && m.cfg.OnDirect != "" {
befehl = m.cfg.OnDirect
}
m.wecken(befehl, map[string]string{
"MM_CHANNEL": kanalID,
"MM_CHANNEL_NAME": kanalName,
"MM_CHANNEL_TYPE": kanalTyp,
"MM_SENDER": autor,
"MM_POST_ID": postID,
"MM_ROOT_ID": wurzel,
"MM_MESSAGE": nachricht,
"MM_TRIGGER": "event",
})
}
// wecken führt den Befehl aus — höchstens einer gleichzeitig, und nicht öfter
// als der Debounce erlaubt. Ohne beides löst ein Gespräch mit fünf Beiträgen
// fünf Agentenläufe aus, die sich gegenseitig überholen.
func (m *Melder) wecken(befehl string, umgebung map[string]string) {
if !m.laufMu.TryLock() {
log.Print("Ein Lauf arbeitet noch — dieser Beitrag wird von ihm mitgenommen.")
return
}
defer m.laufMu.Unlock()
if seit := time.Since(m.letzter); seit < time.Duration(m.cfg.DebounceSeconds)*time.Second {
time.Sleep(time.Duration(m.cfg.DebounceSeconds)*time.Second - seit)
}
m.letzter = time.Now()
log.Printf("Wecke wegen Beitrag von %s in %s", umgebung["MM_SENDER"], umgebung["MM_CHANNEL_NAME"])
cmd := exec.Command("/bin/sh", "-c", befehl)
cmd.Env = os.Environ()
for k, v := range umgebung {
cmd.Env = append(cmd.Env, k+"="+v)
}
cmd.Stdout, cmd.Stderr = os.Stdout, os.Stderr
if err := cmd.Run(); err != nil {
log.Printf("Befehl endete mit Fehler: %v", err)
}
}
// boden sieht auch dann nach, wenn der Stream schweigt. Er holt keine Beiträge,
// sondern führt denselben Befehl aus: Der Agent prüft ohnehin selbst, was
// liegen geblieben ist.
func (m *Melder) boden() {
t := time.NewTicker(time.Duration(m.cfg.FloorSeconds) * time.Second)
defer t.Stop()
for range t.C {
if time.Since(m.letzter) < time.Duration(m.cfg.FloorSeconds)*time.Second {
continue
}
log.Print("Boden: seit Längerem kein Ereignis — sehe trotzdem nach.")
m.wecken(m.cfg.OnMention, map[string]string{"MM_TRIGGER": "floor"})
}
}
func (m *Melder) ichErmitteln() error {
req, _ := http.NewRequest("GET", m.cfg.URL+"/api/v4/users/me", nil)
req.Header.Set("Authorization", "Bearer "+m.cfg.Token)
r, err := http.DefaultClient.Do(req)
if err != nil {
return err
}
defer r.Body.Close()
if r.StatusCode != 200 {
return fmt.Errorf("HTTP %d von /users/me — Token gültig?", r.StatusCode)
}
var u struct {
ID string `json:"id"`
Username string `json:"username"`
}
if err := json.NewDecoder(r.Body).Decode(&u); err != nil {
return err
}
m.ich, m.ichName = u.ID, u.Username
return nil
}