// 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 Mattermost 11.7.10: // // Direktnachricht -> `posted` kommt SOFORT an // eigener Kanalbeitrag -> `posted` kommt an // Kanalbeitrag eines anderen -> kommt NIE an // // Das Gespräch mit einem Agenten läuft also in Echtzeit: Eine DM weckt ihn // ohne Verzögerung. Nur Erwähnungen in Kanälen kommen über den Boden. // // Die Verbindung ist also voll funktionsfähig und hat die Kanal-Abos; es // fehlen ausschließlich die Ereignisse fremder Konten. Vier Varianten wurden // ausgeschlossen: Anmeldung per `authentication_challenge`, Token im // Authorization-Header beim Handshake, ein echter Bot-Account (statt Nutzer // mit Token) und das Sitzungs-Cookie MMAUTHTOKEN wie im Browser. Alle vier // verhalten sich gleich. // // Das Verhalten ist bekannt und nicht auf uns beschränkt: // https://forum.mattermost.com/t/websocket-does-not-provide-posted-events-from-other-accounts/14876 // Dort tritt es auch mit einem System-Admin-Token auf, und eine Lösung steht // nicht dabei. // // Konsequenz für den Betrieb: floor_seconds ist der Weg, nicht der Notnagel, // und der Standard ist deshalb bewusst kurz (15 s). Ein Durchgang kostet einen // API-Aufruf — das ist billig genug, um es oft zu tun, und 15 Sekunden fühlen // sich noch wie eine Antwort an. Wer den Stream repariert bekommt, kann den // Boden hochsetzen; bis dahin trägt er. // // 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 15 — kurz, weil der Stream fremde Beiträge // nicht liefert (siehe Kopf dieser Datei). Auf 0 zu setzen heißt, sich allein // auf den Stream zu verlassen; das hört heute nur die eigenen Beiträge. 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"` // Nach wie vielen Sekunden ohne Arbeit der Punkt von grün auf abwesend // zurückfällt. Standard 120, 0 schaltet die Statusanzeige ganz ab. AwayAfterSeconds int `json:"away_after_seconds"` } 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 marken map[string]int64 // Kanal -> Zeit des zuletzt gesehenen Beitrags markenMu sync.Mutex statusMu sync.Mutex status string // zuletzt gesetzter Status; nur über statusSetzen anfassen aktiv time.Time } 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 = 15 } if cfg.AwayAfterSeconds == 0 { cfg.AwayAfterSeconds = 120 } m := &Melder{cfg: cfg, namen: map[string]string{}, gesehen: map[string]bool{}, marken: map[string]int64{}} 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.") // Offline sagen, statt einen grünen Punkt zurückzulassen, hinter dem // niemand mehr sitzt. m.statusSetzen("offline") os.Exit(0) }() if cfg.FloorSeconds > 0 { go m.boden() } if cfg.AwayAfterSeconds > 0 { go m.abwesendWennStill() } 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 } } }() // Da, aber untätig. Grün ab der ersten Sekunde würde nichts aussagen. m.statusSetzen("away") 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() // Grün heißt: Es passiert gerade etwas, das jemanden angeht. Das routinemäßige // Nachsehen des Bodens gehört ausdrücklich NICHT dazu -- läuft er häufiger als // die Away-Frist, leuchtete der Agent sonst rund um die Uhr, und ein Punkt, // der immer grün ist, sagt genauso wenig wie gar keiner (beim Testen am // 2026-09-01 genau so passiert: Boden alle 20 s, Frist 25 s, nie abwesend). if umgebung["MM_TRIGGER"] != "floor" { m.aktiv = time.Now() m.statusSetzen("online") } 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 nach, wenn der Stream schweigt — aber er weckt NICHT blind. // // Das war der teuerste Fehler dieses Werkzeugs: Der Boden rief den Befehl // einfach alle floor_seconds auf, und weil dahinter ein Sprachmodell hängt, // liefen bei vier Agenten und 15 Sekunden Takt 16 Modellläufe pro Minute -- // rund um die Uhr, auch wenn niemand geschrieben hatte (Andreas, 2026-09-01: // "wachen da permanent die llms auf? das wird dann teuer"). // // goimapnotify hat dieses Problem nicht, weil IMAP IDLE ein echtes Ereignis // liefert: Es feuert, wenn eine Mail kommt, und sonst nie. Genau das gehört // hierher und nicht in den aufgerufenen Befehl -- ein Notifier, der den // Empfänger entscheiden lässt, ob der Anruf nötig war, hat schon angerufen. // // Der Türsteher kostet zwei Anfragen: /users//channels nennt in EINEM // Aufruf alle Kanäle samt last_post_at. Ist keiner neuer als die gemerkte // Marke, passiert nichts. Nur für die geänderten wird der jüngste Beitrag // geholt, um den Autor zu prüfen -- sonst würde die eigene Antwort den // nächsten Durchgang auslösen. // // Gelesen wird der Beitrag NICHT: Der Türsteher fragt "gibt es etwas?", nicht // "was steht drin". Das bleibt beim Agenten, sonst gäbe es zwei Leser und // zwei Antworten. func (m *Melder) boden() { m.markenSetzen() 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 } n, err := m.gibtEsNeues() if err != nil { log.Printf("Boden: Nachsehen fehlgeschlagen: %v", err) continue } if n == "" { continue // still, und das ist die Wahrheit -- niemand wird geweckt } log.Printf("Boden: neuer Beitrag in %s — wecke.", n) m.wecken(m.cfg.OnMention, map[string]string{"MM_TRIGGER": "floor", "MM_CHANNEL": n}) } } // markenSetzen holt den Ausgangsstand: ab wann gilt ein Beitrag als neu. // // Bewusst NICHT über last_post_at aus /users//channels -- das Feld wird auf // dieser Instanz nicht aktualisiert (2026-09-01 gemessen: ein neuer Beitrag // ließ es unverändert). Es ist das dritte Mattermost-Feld nach msg_count und // last_viewed_at, auf das hier kein Verlass ist. Verlässlich ist einzig // `posts?since=`: leer, wenn nichts Neues da ist. func (m *Melder) markenSetzen() { kanaele, err := m.meineKanaele() if err != nil { log.Printf("Ausgangsstand nicht ermittelbar: %v", err) return } // Ein kurzer Rückgriff, damit nicht verlorengeht, was unmittelbar vor dem // Start eintraf -- aber kein langer, sonst weckt der erste Durchgang für // alles, was heute geschrieben wurde. start := time.Now().Add(-2 * time.Minute).UnixMilli() m.markenMu.Lock() for _, k := range kanaele { m.marken[k.ID] = start } m.markenMu.Unlock() log.Printf("Ausgangsstand über %d Kanäle gesetzt.", len(kanaele)) } type kanal struct { ID string `json:"id"` Name string `json:"name"` Type string `json:"type"` LastPostAt int64 `json:"last_post_at"` } func (m *Melder) meineKanaele() ([]kanal, error) { req, _ := http.NewRequest("GET", m.cfg.URL+"/api/v4/users/"+m.ich+"/channels", nil) req.Header.Set("Authorization", "Bearer "+m.cfg.Token) r, err := http.DefaultClient.Do(req) if err != nil { return nil, err } defer r.Body.Close() if r.StatusCode != 200 { return nil, fmt.Errorf("HTTP %d", r.StatusCode) } var ks []kanal return ks, json.NewDecoder(r.Body).Decode(&ks) } // gibtEsNeues nennt den ersten Kanal mit einem fremden neuen Beitrag, sonst "". // // Ein Aufruf je Kanal, aber ein billiger: `posts?since=` antwortet bei // Stille mit einer leeren Liste. Bei einem Dutzend Kanälen und einem Takt von // 60 s sind das zwölf Anfragen pro Minute -- gegen einen Modelllauf, der sonst // jedes Mal liefe, ist das nichts. // // Gelesen wird der Beitrag nicht: Der Türsteher fragt "gibt es etwas?", nicht // "was steht drin". Das bleibt beim Agenten -- zwei Leser wären zwei Antworten. func (m *Melder) gibtEsNeues() (string, error) { kanaele, err := m.meineKanaele() if err != nil { return "", err } jetzt := time.Now().UnixMilli() for _, k := range kanaele { m.markenMu.Lock() marke, bekannt := m.marken[k.ID] if !bekannt { // Neuer Kanal: ab jetzt beobachten, nicht rückwirkend wecken. m.marken[k.ID] = jetzt m.markenMu.Unlock() continue } m.markenMu.Unlock() fremd, err := m.neuesSeit(k.ID, marke) if err != nil { continue // ein Kanal, der klemmt, darf die anderen nicht aufhalten } m.markenMu.Lock() m.marken[k.ID] = jetzt m.markenMu.Unlock() if fremd { if k.Name != "" { return k.Name, nil } return k.ID, nil } } return "", nil } // neuesSeit sagt, ob seit `marke` ein Beitrag von jemand anderem kam. Die // eigene Antwort zählt nicht -- sonst löste sie den nächsten Durchgang aus. func (m *Melder) neuesSeit(kanalID string, marke int64) (bool, error) { req, _ := http.NewRequest("GET", fmt.Sprintf("%s/api/v4/channels/%s/posts?since=%d", m.cfg.URL, kanalID, marke), nil) req.Header.Set("Authorization", "Bearer "+m.cfg.Token) r, err := http.DefaultClient.Do(req) if err != nil { return false, err } defer r.Body.Close() if r.StatusCode != 200 { return false, fmt.Errorf("HTTP %d", r.StatusCode) } var a struct { Order []string `json:"order"` Posts map[string]map[string]any `json:"posts"` } if err := json.NewDecoder(r.Body).Decode(&a); err != nil { return false, err } for _, id := range a.Order { p := a.Posts[id] if uid, _ := p["user_id"].(string); uid == m.ich { continue } if typ, _ := p["type"].(string); strings.HasPrefix(typ, "system_") { continue } return true, nil } return false, nil } // status sagt grün oder abwesend — und darf dabei nie etwas kaputtmachen. // // Der Punkt neben dem Namen ist eine Höflichkeit gegenüber dem, der in die // Seitenleiste schaut. Weigert sich der Server, muss die Arbeit trotzdem // stattfinden: Ein Agent, der aufhört zu antworten, weil er einen Punkt nicht // einfärben konnte, wäre absurd (übernommen aus hive, channels/team.py). // // Warum überhaupt zwei Zustände: Ein Agent ganz ohne Status sieht aus, als sei // niemand da. Einer, der rund um die Uhr grün leuchtet, sagt genauso wenig. // Beides ist schlechter als die Wahrheit — da, und das ist der Zeitpunkt, an // dem zuletzt etwas geschah. func (m *Melder) statusSetzen(neu string) { m.statusMu.Lock() if m.cfg.AwayAfterSeconds < 0 || m.status == neu { m.statusMu.Unlock() return } m.statusMu.Unlock() koerper, _ := json.Marshal(map[string]string{"user_id": m.ich, "status": neu}) req, _ := http.NewRequest("PUT", m.cfg.URL+"/api/v4/users/"+m.ich+"/status", strings.NewReader(string(koerper))) req.Header.Set("Authorization", "Bearer "+m.cfg.Token) req.Header.Set("Content-Type", "application/json") r, err := http.DefaultClient.Do(req) if err != nil { log.Printf("Status %q nicht gesetzt: %v", neu, err) return } defer r.Body.Close() if r.StatusCode != 200 { log.Printf("Status %q nicht gesetzt: HTTP %d", neu, r.StatusCode) return } m.statusMu.Lock() m.status = neu m.statusMu.Unlock() } // abwesendWennStill lässt das Grün verblassen, sobald es still geworden ist. func (m *Melder) abwesendWennStill() { t := time.NewTicker(15 * time.Second) defer t.Stop() for range t.C { if !m.aktiv.IsZero() && time.Since(m.aktiv) > time.Duration(m.cfg.AwayAfterSeconds)*time.Second { m.statusSetzen("away") } } } 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 }