Files
mmnotify/main.go
T
elton 831c2a953f Direktnachrichten kommen doch sofort an -- das Gespraech ist Echtzeit
Andreas: "dann durchsuch noch mal die foren, das kann doch nicht sein
das das niemand da draussen hat das problem." Hatte jemand, und der
entscheidende Satz stand in einem Bugreport eines anderen Projekts:
"DMs work fine."

Nachgemessen, und es stimmt:

  Direktnachricht            -> posted kommt SOFORT
  eigener Kanalbeitrag       -> posted kommt
  Kanalbeitrag eines anderen -> kommt nie

Damit ist die Lage eine voellig andere als gedacht. Mit einem Agenten zu
reden -- der Anwendungsfall, um den es Andreas ging -- laeuft in Echtzeit
ueber eine DM. Der Boden deckt nur noch Erwaehnungen in Kanaelen ab, und
das ist der Fall, bei dem 15 Sekunden niemanden stoeren.

Dabei ausserdem einen eigenen Fehler korrigiert: Ich hatte den Test
"direkt am Container, ohne Caddy" als erledigt dargestellt, obwohl das
sed dafuer fehlgeschlagen war und gegen den Proxy gemessen wurde.
Nachgeholt: ws://192.168.1.142:8065 direkt verhaelt sich genauso. Der
Proxy ist jetzt wirklich ausgeschlossen.
2026-09-01 20:54:02 +02:00

445 lines
14 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 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
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{}}
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 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"})
}
}
// 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
}