Files
mesh-catalog/modules/messenger/cmd/messenger/main.go
T
jochen 0ae7933d54 Tell the operator what the mesh finds wrong, and watch the watcher (hq to-be 45 phase 1)
The mesh noticed 48 core failures in six days and told nobody (ADR 0227).
messenger holds the operator-channel seat: it consumes the controller's
condition events and sends them to Telegram and the desktop notifier,
deduplicated by key, reminded once, edited on clear, capped at 20 an hour
with the rest folded, and refusing anything carrying an address, a path or
a secret. mesh-watcher, on a machine other than the control node, sends to
Telegram directly when the self-check heartbeat or the bus goes silent.
2026-10-06 09:35:55 +02:00

300 lines
11 KiB
Go

// messenger: the holder of the operator-channel seat (novox/hq to-be 45 §5, ADR 0227, research 028).
// A Go bundle the node's runtime launches. It consumes the controller's condition events and decides
// what is said to the operator, on Telegram and on the desktop notifier of the machine the operator
// is at. It keeps its open messages in its own state, so a restart forgets nothing. stdout is the MCP
// channel; what this module says, it says on stderr.
package main
import (
"encoding/json"
"errors"
"fmt"
"os"
"strings"
"sync"
"time"
stdio "git.novox.be/novox/mesh-sdk/go"
)
func errorf(format string, a ...any) error { return fmt.Errorf(format, a...) }
func logf(format string, a ...any) { fmt.Fprintf(os.Stderr, format+"\n", a...) }
// Settings are the operator's values for this module, merged by the mesh into one JSON file.
type Settings struct {
TelegramChatID string `json:"telegram-chat-id"`
DesktopMachines []string `json:"desktop-machines"`
}
func readSettings(path string) (Settings, error) {
var s Settings
if path == "" {
return s, errors.New("started without a settings file")
}
raw, err := os.ReadFile(path)
if err != nil {
return s, err
}
// The chat id may be given as a number.
var loose map[string]any
if err := json.Unmarshal(raw, &loose); err != nil {
return s, fmt.Errorf("the settings file is not JSON: %v", err)
}
switch v := loose["telegram-chat-id"].(type) {
case string:
s.TelegramChatID = strings.TrimSpace(v)
case float64:
s.TelegramChatID = fmt.Sprintf("%.0f", v)
}
if list, ok := loose["desktop-machines"].([]any); ok {
for _, m := range list {
if name, ok := m.(string); ok && strings.TrimSpace(name) != "" {
s.DesktopMachines = append(s.DesktopMachines, strings.TrimSpace(name))
}
}
}
return s, nil
}
// stateStore keeps the open messages in the module's declared state `open`, and the recent sends
// under one key of `sent` (ADR 0201).
type stateStore struct{}
// kvKey is a condition key as a bucket key: only the characters a key may carry.
func kvKey(key string) string {
var b strings.Builder
for _, r := range key {
switch {
case r >= 'a' && r <= 'z', r >= 'A' && r <= 'Z', r >= '0' && r <= '9', r == '-', r == '_', r == '.', r == '=':
b.WriteRune(r)
default:
b.WriteRune('_')
}
}
return strings.Trim(b.String(), ".")
}
func (stateStore) Put(r Record) error {
_, err := stdio.State("open").Put(kvKey(r.Key), r)
return err
}
func (stateStore) Delete(key string) error { return stdio.State("open").Delete(kvKey(key)) }
func (stateStore) All() ([]Record, error) {
keys, err := stdio.State("open").Keys()
if err != nil {
return nil, err
}
var out []Record
for _, k := range keys {
e, err := stdio.State("open").Get(k)
if err != nil {
return nil, err
}
if e == nil {
continue
}
var r Record
if err := json.Unmarshal(e.Value, &r); err != nil {
logf("[messenger] the open message kept as %s cannot be read (%v); left as it is", k, err)
continue
}
out = append(out, r)
}
return out, nil
}
func (stateStore) PutRecent(s []Sent) error {
_, err := stdio.State("sent").Put("recent", s)
return err
}
func (stateStore) Recent() ([]Sent, error) {
e, err := stdio.State("sent").Get("recent")
if err != nil || e == nil {
return nil, err
}
var s []Sent
return s, json.Unmarshal(e.Value, &s)
}
// listening is whether condition events reach this holder, in words.
type listening struct {
mu sync.Mutex
now string
}
func (l *listening) set(s string) { l.mu.Lock(); l.now = s; l.mu.Unlock() }
func (l *listening) get() string { l.mu.Lock(); defer l.mu.Unlock(); return l.now }
func main() {
settingsFile := os.Getenv("MESH_MESSENGER_SETTINGS")
var said sync.Mutex
lastSaid := ""
settings := func() Settings {
s, err := readSettings(settingsFile)
said.Lock()
defer said.Unlock()
if err != nil && err.Error() != lastSaid {
logf("[messenger] settings cannot be read: %v", err)
}
lastSaid = ""
if err != nil {
lastSaid = err.Error()
}
return s
}
h := &Holder{
Telegram: NewTelegram(TelegramConfig{
TokenFile: os.Getenv("MESH_MESSENGER_TELEGRAM_TOKEN_FILE"),
ChatID: func() string { return settings().TelegramChatID },
}),
Desktop: &Desktop{
Machines: func() []string { return settings().DesktopMachines },
Ask: stdio.Ask,
},
Store: stateStore{},
Logf: logf,
Emit: func(event string, body any) error { return stdio.Emit(event, body) },
}
h.init()
l := &listening{now: "not yet: starting"}
go run(h, l)
if err := stdio.Serve("", tools(h, l)); err != nil {
logf("%v", err)
os.Exit(1)
}
}
// run reads back the state, listens for condition events, and keeps time — each retried, and each
// failure said, never given up on quietly.
func run(h *Holder, l *listening) {
time.Sleep(500 * time.Millisecond) // Serve first: the state is reached through it
for wait := 2 * time.Second; ; wait = min(wait*2, time.Minute) {
err := h.Load()
if err == nil {
break
}
logf("[messenger] cannot read back its open messages yet (%v); asking again in %s", err, wait)
time.Sleep(wait)
}
for _, ch := range h.channels() {
if err := ch.Ready(); err != nil {
logf("[messenger] %s cannot send: %v", ch.Name(), err)
}
}
go func() {
for range time.Tick(time.Minute) {
h.Tick()
}
}()
handle := func(e stdio.Envelope) error {
event, ok := eventOf(e.Key)
if !ok {
return nil
}
c, err := DecodeCondition(event, e.Body)
if err != nil {
h.Unreadable(e.Key, err)
return nil
}
h.Condition(event, c)
return nil
}
for wait := 2 * time.Second; ; wait = min(wait*2, time.Minute) {
err := stdio.Subscribe(ControllerSeat+".*", handle)
if err == nil {
l.set("listening")
logf("[messenger] listening for the controller's condition events")
return
}
l.set("not yet: " + err.Error())
logf("[messenger] not hearing condition events yet (%v); asking again in %s", err, wait)
time.Sleep(wait)
}
}
func str(description string) map[string]any {
return map[string]any{"type": "string", "description": description}
}
func strArg(a map[string]any, k string) string { s, _ := a[k].(string); return strings.TrimSpace(s) }
func limitArg(a map[string]any, def int) int {
if v, ok := a["limit"].(float64); ok && v >= 1 {
return min(int(v), KeptSends)
}
return def
}
func tools(h *Holder, l *listening) []stdio.Tool {
return []stdio.Tool{
{Name: "operator-channel.open",
Description: "What is open now: every message the operator was sent about something still wrong, urgent first, " +
"oldest first — its key, severity, summary, since when, the channels it went to, whether it is silenced, " +
"reminded, held by the cap, refused, or not sent yet.",
Run: func(map[string]any) (any, error) { return h.Open(), nil }},
{Name: "operator-channel.history",
Description: "What was said to the operator lately, newest first — each send, edit, fold and failure with its " +
"channel, key and outcome — and every message refused for carrying an address, a path or a secret.",
Input: map[string]any{"limit": map[string]any{"type": "integer", "description": "at most this many of each (default 50)"}},
Run: func(a map[string]any) (any, error) { return h.History(limitArg(a, 50)), nil }},
{Name: "operator-channel.notify",
Description: "Tell the operator something, as a module that uses the seat: a key (the same key is the same " +
"message), urgent or warning, one line in the mesh's words, and what it is about. clear says it is over. " +
"Roles and words only: an address, a path or a secret is refused. Deduplicated, capped and routed " +
"like the controller's conditions.",
Input: map[string]any{
"key": str("what makes it the same message the next time, e.g. backup.ace.failed"),
"severity": map[string]any{"type": "string", "enum": []string{Urgent, Warning}},
"summary": str("one line in the mesh's words"),
"subject": str("what it is about: a machine's role, a module, a plan"),
"clear": map[string]any{"type": "boolean", "description": "it is over: the message is edited to say so"},
},
Run: func(a map[string]any) (any, error) {
clear, _ := a["clear"].(bool)
return h.Notify(strArg(a, "key"), strArg(a, "severity"), strArg(a, "summary"), strArg(a, "subject"), clear)
}},
{Name: "messenger_status",
Description: "Whether the operator can be reached, and why not: each channel — can it send, what it lacks " +
"(the Telegram bot token and chat id, the desktop machines), its last error, how many it sent in the " +
"last hour against the cap, how many it holds — the open and silenced count, messages not sent yet, " +
"refusals, unreadable events, whether condition events arrive, and the rules it applies. check asks " +
"Telegram who the bot is, sending nothing.",
Input: map[string]any{"check": map[string]any{"type": "boolean", "description": "also ask Telegram whether the token works"}},
Run: func(a map[string]any) (any, error) {
st := h.Status(l.get())
if check, _ := a["check"].(bool); check {
if t, ok := h.Telegram.(*Telegram); ok {
who, err := t.Who()
if err != nil {
return map[string]any{"status": st, "telegram_check": "failed: " + err.Error()}, nil
}
return map[string]any{"status": st, "telegram_check": "the token works: the bot is @" + who}, nil
}
}
return st, nil
}},
{Name: "messenger_recent",
Description: "The recent sends — each message, edit, fold and failure on each channel, newest first — and the refusals.",
Input: map[string]any{"limit": map[string]any{"type": "integer", "description": "at most this many (default 20)"}},
Run: func(a map[string]any) (any, error) { return h.History(limitArg(a, 20)), nil }},
{Name: "messenger_test",
Description: "Send a test message now, to telegram, desktop or both (default both), through the content rule " +
"and the cap: proves a channel reaches the operator. Answers per channel: sent, or why not.",
Input: map[string]any{
"channel": map[string]any{"type": "string", "enum": []string{"telegram", "desktop", "both"}},
"text": str("the words (default: a line saying it is a test)"),
},
Run: func(a map[string]any) (any, error) { return h.Test(strArg(a, "channel"), strArg(a, "text")), nil }},
{Name: "messenger_check",
Description: "Would these words be allowed to leave the mesh? Runs the content rule — no address, path or " +
"secret — and answers allowed, or the class of what it carried. Sends nothing.",
Input: map[string]any{"text": str("the words to judge")},
Run: func(a map[string]any) (any, error) {
if r, ok := Check(strArg(a, "text")); !ok {
return map[string]any{"allowed": false, "class": r.Class, "carried": r.What}, nil
}
return map[string]any{"allowed": true}, nil
}},
}
}