messenger: never say history as news — read the open conditions first, coalesce bursts
A consumer made on 2026-10-06 was handed three hours of raised-and-cleared conditions at once and the holder said each as new: 20 desktop notifications in a second. What is said is now decided by the controller's open set and by an event's own time, never by its arrival; bursts are one message, the desktop gets warnings at most every 15 min, the cap is said once, and the first minute after start says only the urgent conditions still open. Replays of that morning's 96 events are tests. Also drops two committed binaries. (novox/hq issue 271)
This commit is contained in:
@@ -6,3 +6,6 @@ modules/slack/cmd/slack-tools/slack-tools
|
||||
modules/jetbrains-toolbox/cmd/toolbox-tools/toolbox-tools
|
||||
modules/messenger/cmd/messenger/messenger
|
||||
modules/mesh-watcher/cmd/mesh-watcher/mesh-watcher
|
||||
# ...and the same built from the module's root (go build ./cmd/<name>), which names it after the module.
|
||||
modules/messenger/messenger
|
||||
modules/mesh-watcher/mesh-watcher
|
||||
|
||||
Binary file not shown.
@@ -5,6 +5,8 @@ to-be 45 §5, ADR 0227, research 028 — the minimal form, Q1a, Q5a, Q8).
|
||||
|
||||
- **Declares and holds `operator-channel`**, held once for the mesh, serving `open`, `history` and
|
||||
`notify`.
|
||||
- **Reads the controller's open conditions** (`seat:mesh-controller.conditions`) — the state it
|
||||
believes, before and beside the events.
|
||||
- **Consumes the controller's condition events** — `mesh-controller.condition-raised`, `-changed`,
|
||||
`-cleared` — and decides what is sent. The controller calls nobody. The events are read in one
|
||||
place, `cmd/messenger/condition.go`, which states every assumption it makes about their shape.
|
||||
@@ -16,26 +18,51 @@ to-be 45 §5, ADR 0227, research 028 — the minimal form, Q1a, Q5a, Q8).
|
||||
|
||||
## What is sent, and when
|
||||
|
||||
**History is never news** (novox/hq issue 271). The events stream hands everything it holds to a
|
||||
consumer that is new, was made again, or whose holder was away — on 2026-10-06 a newly made consumer
|
||||
handed over three hours of raised-and-cleared conditions at once, and each was said as if new. So
|
||||
what is said is decided by the state now and by when a thing happened, never by when it arrived:
|
||||
|
||||
- **The state first.** On start, every 10 min, and at once when an old event arrives, the holder reads
|
||||
what is open now from the controller (`mesh-controller.conditions`). An event older than that
|
||||
reading is already in it and is not acted on. What is open here and not there ended meanwhile: it is
|
||||
forgotten quietly (Telegram's message is edited, which notifies nobody; the desktop is not woken). An
|
||||
urgent condition open and never said is said — all of them in **one** message. A warning raised in
|
||||
the last 10 min and not heard is said like any other; an older one is state.
|
||||
- **Old is state.** An event whose own time (`at`) is more than **10 min** old is recorded, never said.
|
||||
- **A clearing of something never said says nothing**, and a message still held when its condition
|
||||
clears is dropped.
|
||||
- **Start grace.** Nothing is sent in the **first minute** after start but that one urgent message;
|
||||
what arrives meanwhile goes out, coalesced, when the minute ends.
|
||||
|
||||
| When | What |
|
||||
|---|---|
|
||||
| `condition-raised` | one message, deduplicated by the condition's key |
|
||||
| still open after 1 h (urgent) or 12 h (warning) | once more |
|
||||
| still open after 1 h (urgent) or 12 h (warning) | once more — never sooner than that after it was first said |
|
||||
| `condition-changed` from warning to urgent | once more, to both channels |
|
||||
| `condition-cleared` | the first message edited to say so (both channels can) |
|
||||
| `condition-cleared` | the first message edited to say so; said inside a digest, a line in the next one |
|
||||
| cleared and raised again within 10 min | the same message, edited back to open — not a new one |
|
||||
| silenced | nothing, its clearing included |
|
||||
|
||||
- **Routing:** urgent to Telegram and the desktop; warning to the desktop when a session there
|
||||
answers, otherwise to Telegram. The notifier answering is how "the operator's session is there" is
|
||||
read until presence is decided (research 028 Q4).
|
||||
- **Rate:** at most 20 messages an hour per channel. The rest are held and folded into one message
|
||||
naming them all, sent at most every ten minutes — the cap is said, never silent.
|
||||
- **Bursts are one message.** What is to be said on a channel is held **30 s** from the first and goes
|
||||
out together: alone, as itself; several, as one digest (`WARNING: 5 new warnings`, a line each). An
|
||||
urgent message goes at once when nothing went out on its channel in the last 30 s.
|
||||
- **The desktop is gentle:** warnings reach it as at most **one message every 15 min**; urgent ones
|
||||
are not held to that.
|
||||
- **Cap:** at most 20 messages an hour per channel — a last line of defence the above keeps far away.
|
||||
Reached, it is said **once**, in a message of its own, and the rest is held to go out as one digest
|
||||
when the hour allows.
|
||||
- **What may leave the mesh:** roles and words. A message carrying an address (IP, host name, URL,
|
||||
mail address), a path, or anything shaped like a secret (a token, a key block, a long random or
|
||||
hexadecimal string, `password=…`) is refused, logged, stated as the `refused` event, and sent in its
|
||||
place as `channel-refused` with the words that carried it withheld. The rule is `cmd/messenger/content.go`.
|
||||
- **Nothing silent:** a channel that cannot send says so in `messenger_status` and in the log, and
|
||||
the message is tried again every minute while the condition is open. An event that cannot be read
|
||||
what it holds is tried again every minute while the condition is open. `messenger_status` also says
|
||||
when the open conditions were last read, why they could not be, and how many events were taken as
|
||||
history. An event that cannot be read
|
||||
is refused by name, counted, and told to the operator once an hour.
|
||||
- **No answering back.** Acknowledging is `conditions silence`, through the mesh.
|
||||
|
||||
@@ -58,7 +85,7 @@ Nothing is sent until these are given; `messenger_status` says which is missing.
|
||||
| `operator-channel.history` | what was said lately, and the refusals |
|
||||
| `operator-channel.notify` | a message from a module using the seat: key, severity, summary; `clear` to end it |
|
||||
| `messenger_status` | whether the operator can be reached and why not; `check` asks Telegram whether the token works |
|
||||
| `messenger_recent` | the recent sends, edits, folds and failures |
|
||||
| `messenger_recent` | the recent sends, edits, digests and failures |
|
||||
| `messenger_test` | a test message now, to telegram, desktop or both |
|
||||
| `messenger_check` | whether some words may leave the mesh |
|
||||
|
||||
|
||||
@@ -20,6 +20,11 @@ package main
|
||||
// - `key` is `<scope>.<id>.<kind>`. When it is absent it is made from subject and kind; when
|
||||
// neither gives one the event is unreadable.
|
||||
// - A cleared event carries the condition as last held, and may add `cleared` (its time).
|
||||
// - Every event carries `at`, when the transition happened (the controller's conditions.Event).
|
||||
// It is how this holder tells news from history: a stream replays what it holds to a consumer
|
||||
// that is new or was away, and an event's delivery time says nothing about when it happened.
|
||||
// - The `conditions` verb of the controller's seat answers `{conditions: [...], open, note}`, each
|
||||
// entry the condition as the store holds it — the same shape as an event's body.
|
||||
//
|
||||
// An event this cannot read is refused by name — which event, which field, why — counted, said in
|
||||
// the status and in the log, and told to the operator; never read as an empty condition.
|
||||
@@ -63,6 +68,22 @@ type Condition struct {
|
||||
SilencedTill time.Time
|
||||
SilencedWhy string
|
||||
Cleared time.Time
|
||||
// At is when the transition the event says happened; zero for a condition read from the open set.
|
||||
At time.Time
|
||||
}
|
||||
|
||||
// When is the event's own time: `at`; without it, the clearing's or the raising's time; without
|
||||
// those, zero — and a zero time is read as now by the holder, never as old.
|
||||
func (c Condition) When(event string) time.Time {
|
||||
switch {
|
||||
case !c.At.IsZero():
|
||||
return c.At
|
||||
case event == EventCleared:
|
||||
return c.Cleared
|
||||
case event == EventRaised:
|
||||
return c.Raised
|
||||
}
|
||||
return time.Time{}
|
||||
}
|
||||
|
||||
// SubjectWords is the subject as a person reads it: "machine ace", "plan 41", "provider keycloak".
|
||||
@@ -147,6 +168,7 @@ func DecodeCondition(event string, body []byte) (Condition, error) {
|
||||
c.Raised = when("raised")
|
||||
c.LastObserved = when("last-observed", "last_observed", "lastObserved")
|
||||
c.Cleared = when("cleared")
|
||||
c.At = when("at")
|
||||
if err != nil {
|
||||
return Condition{}, err
|
||||
}
|
||||
@@ -249,3 +271,55 @@ func DecodeCondition(event string, body []byte) (Condition, error) {
|
||||
}
|
||||
return c, nil
|
||||
}
|
||||
|
||||
// DecodeOpen reads the controller's `conditions` answer — what is open now — bare, inside the
|
||||
// runtime's `{answer}` or `{output}`, inside a tool reply's text content, or as a JSON string. An
|
||||
// answer without a `conditions` list is refused: an empty open set is said, never assumed.
|
||||
func DecodeOpen(raw json.RawMessage) ([]Condition, error) {
|
||||
var m map[string]json.RawMessage
|
||||
if err := json.Unmarshal(raw, &m); err != nil {
|
||||
var s string
|
||||
if json.Unmarshal(raw, &s) == nil {
|
||||
return DecodeOpen(json.RawMessage(s))
|
||||
}
|
||||
return nil, fmt.Errorf("the conditions answer is not JSON")
|
||||
}
|
||||
if list, ok := m["conditions"]; ok {
|
||||
var items []json.RawMessage
|
||||
if string(list) != "null" {
|
||||
if err := json.Unmarshal(list, &items); err != nil {
|
||||
return nil, fmt.Errorf("the conditions answer's conditions is not a list")
|
||||
}
|
||||
}
|
||||
out := make([]Condition, 0, len(items))
|
||||
for _, it := range items {
|
||||
c, err := DecodeCondition("conditions", it)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
out = append(out, c)
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
for _, k := range []string{"answer", "output"} {
|
||||
if v, ok := m[k]; ok && string(v) != "null" {
|
||||
return DecodeOpen(v)
|
||||
}
|
||||
}
|
||||
if v, ok := m["content"]; ok {
|
||||
var content []struct {
|
||||
Text string `json:"text"`
|
||||
}
|
||||
var e struct {
|
||||
IsError bool `json:"isError"`
|
||||
}
|
||||
_ = json.Unmarshal(raw, &e)
|
||||
if json.Unmarshal(v, &content) == nil && len(content) > 0 {
|
||||
if e.IsError {
|
||||
return nil, fmt.Errorf("the controller refused: %s", content[0].Text)
|
||||
}
|
||||
return DecodeOpen(json.RawMessage(content[0].Text))
|
||||
}
|
||||
}
|
||||
return nil, fmt.Errorf("the conditions answer lists no conditions")
|
||||
}
|
||||
|
||||
@@ -33,6 +33,9 @@ type Desktop struct {
|
||||
func (d *Desktop) Name() string { return "desktop" }
|
||||
func (d *Desktop) CanEdit() bool { return true }
|
||||
|
||||
// SilentEdit: replacing a notification shows it again.
|
||||
func (d *Desktop) SilentEdit() bool { return false }
|
||||
|
||||
func (d *Desktop) Ready() error {
|
||||
if len(d.Machines()) == 0 {
|
||||
return errors.New("no machine to show it on: the setting desktop-machines is not given " +
|
||||
|
||||
@@ -0,0 +1,460 @@
|
||||
package main
|
||||
|
||||
// History is never news (novox/hq issue 271). testdata/backlog-2026-10-06.json is the controller's
|
||||
// own record of every condition event of that morning — 96 of them, raised and cleared between 08:32
|
||||
// and 09:49 UTC — exactly the bodies the events stream held when the holder's consumer was first made
|
||||
// at 09:48, and handed over in one go. The holder of that day sent 20 desktop notifications in one
|
||||
// second and folded the rest. These replay it.
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"os"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
type backlogEvent struct {
|
||||
event string
|
||||
body json.RawMessage
|
||||
at time.Time
|
||||
}
|
||||
|
||||
func loadBacklog(t *testing.T) []backlogEvent {
|
||||
t.Helper()
|
||||
raw, err := os.ReadFile("testdata/backlog-2026-10-06.json")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
var bodies []json.RawMessage
|
||||
if err := json.Unmarshal(raw, &bodies); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
var out []backlogEvent
|
||||
for _, b := range bodies {
|
||||
var head struct {
|
||||
Event string `json:"event"`
|
||||
At time.Time `json:"at"`
|
||||
}
|
||||
if err := json.Unmarshal(b, &head); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
out = append(out, backlogEvent{head.Event, b, head.At})
|
||||
}
|
||||
if len(out) != 96 {
|
||||
t.Fatalf("the backlog has %d events, want the 96 of that morning", len(out))
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
// openAt is the controller's open set at a moment, as the backlog makes it: what was raised and not
|
||||
// cleared by then.
|
||||
func openAt(t *testing.T, backlog []backlogEvent, at time.Time) []Condition {
|
||||
t.Helper()
|
||||
open := map[string]Condition{}
|
||||
var order []string
|
||||
for _, e := range backlog {
|
||||
if e.at.After(at) {
|
||||
continue
|
||||
}
|
||||
c, err := DecodeCondition(e.event, e.body)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
switch e.event {
|
||||
case EventCleared:
|
||||
delete(open, c.Key)
|
||||
default:
|
||||
if _, ok := open[c.Key]; !ok {
|
||||
order = append(order, c.Key)
|
||||
}
|
||||
c.At = time.Time{}
|
||||
open[c.Key] = c
|
||||
}
|
||||
}
|
||||
var out []Condition
|
||||
for _, k := range order {
|
||||
if c, ok := open[k]; ok {
|
||||
out = append(out, c)
|
||||
}
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
// replay hands the holder every event of the backlog the way run() does: read by name, then taken.
|
||||
func replay(t *testing.T, h *Holder, backlog []backlogEvent, skip func(Condition, string) bool) {
|
||||
t.Helper()
|
||||
for _, e := range backlog {
|
||||
event, ok := eventOf(ControllerSeat + "." + e.event)
|
||||
if !ok {
|
||||
t.Fatalf("%s is not a condition event", e.event)
|
||||
}
|
||||
c, err := DecodeCondition(event, e.body)
|
||||
if err != nil {
|
||||
t.Fatalf("today's event is unreadable: %v", err)
|
||||
}
|
||||
if skip != nil && skip(c, event) {
|
||||
continue
|
||||
}
|
||||
h.Condition(event, c)
|
||||
}
|
||||
}
|
||||
|
||||
// runFor lets time pass as main does: a flush every 5 s, a tick every minute.
|
||||
func runFor(h *Holder, c *clock, d time.Duration) {
|
||||
for passed := time.Duration(0); passed < d; passed += 5 * time.Second {
|
||||
c.pass(5 * time.Second)
|
||||
h.Flush()
|
||||
if passed%time.Minute == 0 {
|
||||
h.Tick()
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func said(chs ...*fakeChannel) int {
|
||||
n := 0
|
||||
for _, ch := range chs {
|
||||
n += len(ch.sends) + len(ch.edits)
|
||||
}
|
||||
return n
|
||||
}
|
||||
|
||||
// The consumer was made at 09:48:46 UTC; the holder had started a moment before.
|
||||
func incidentHolder(t *testing.T) (*Holder, *fakeChannel, *fakeChannel, *clock) {
|
||||
h, tg, dt, c, _ := newHolder(t)
|
||||
c.t = time.Date(2026, 10, 6, 9, 49, 30, 0, time.UTC)
|
||||
h.Start()
|
||||
return h, tg, dt, c
|
||||
}
|
||||
|
||||
func TestTodaysBacklogReplayedSaysNothing(t *testing.T) {
|
||||
backlog := loadBacklog(t)
|
||||
h, tg, dt, c := incidentHolder(t)
|
||||
open := openAt(t, backlog, c.now())
|
||||
if len(open) != 0 {
|
||||
t.Fatalf("at 09:49 the controller had %d open; the record says none", len(open))
|
||||
}
|
||||
h.Sync(open, c.now())
|
||||
c.pass(time.Second)
|
||||
replay(t, h, backlog, nil)
|
||||
runFor(h, c, 30*time.Minute)
|
||||
if n := said(tg, dt); n != 0 {
|
||||
t.Fatalf("the replay of that morning said %d thing(s): telegram %v desktop %v", n, tg.sends, dt.sends)
|
||||
}
|
||||
st := h.Status("listening")
|
||||
if st.History != 96 || st.Open != 0 {
|
||||
t.Fatalf("events already in the state %d (want all 96), open %d", st.History, st.Open)
|
||||
}
|
||||
}
|
||||
|
||||
func TestTodaysBacklogWithTheControllerUnreachableSaysNothing(t *testing.T) {
|
||||
backlog := loadBacklog(t)
|
||||
h, tg, dt, c := incidentHolder(t)
|
||||
h.SyncFailed(fmt.Errorf("no answer in 30 s"))
|
||||
replay(t, h, backlog, nil)
|
||||
runFor(h, c, 30*time.Minute)
|
||||
if n := said(tg, dt); n != 0 {
|
||||
t.Fatalf("said %d: telegram %v desktop %v", n, tg.sends, dt.sends)
|
||||
}
|
||||
// Old by their own time, all but the last two: state, and the state is asked for again.
|
||||
if st := h.Status("listening"); st.Stale < 90 || st.Open != 0 || !h.WantsSync() {
|
||||
t.Fatalf("taken as state %d, open %d, wants the state read: %v", st.Stale, st.Open, h.WantsSync())
|
||||
}
|
||||
}
|
||||
|
||||
func TestTodaysBacklogRedeliveredAndReplayedAgainLaterSaysNothing(t *testing.T) {
|
||||
backlog := loadBacklog(t)
|
||||
h, tg, dt, c := incidentHolder(t)
|
||||
h.Sync(openAt(t, backlog, c.now()), c.now())
|
||||
replay(t, h, backlog, nil)
|
||||
replay(t, h, backlog, nil) // redelivered: unacknowledged in time
|
||||
runFor(h, c, time.Hour)
|
||||
// An hour on, the consumer is deleted and made again: the stream hands the morning over once more.
|
||||
h.SyncFailed(fmt.Errorf("the bus is restarting"))
|
||||
replay(t, h, backlog, nil)
|
||||
runFor(h, c, 30*time.Minute)
|
||||
if n := said(tg, dt); n != 0 {
|
||||
t.Fatalf("said %d: telegram %v desktop %v", n, tg.sends, dt.sends)
|
||||
}
|
||||
}
|
||||
|
||||
func TestTodaysBacklogWithAnUrgentStillOpenIsOneMessage(t *testing.T) {
|
||||
backlog := loadBacklog(t)
|
||||
h, tg, dt, c := incidentHolder(t)
|
||||
// Suppose the merge never got acted on, and three machines stayed behind: still open at 09:49.
|
||||
stillOpen := map[string]bool{"merge.mesh-catalog.f2b51686.not-acted": true,
|
||||
"machine.ace.core-behind": true, "machine.g14.core-behind": true, "machine.shanks.core-behind": true}
|
||||
var kept []backlogEvent
|
||||
for _, e := range backlog {
|
||||
c, _ := DecodeCondition(e.event, e.body)
|
||||
if e.event == EventCleared && stillOpen[c.Key] {
|
||||
continue
|
||||
}
|
||||
kept = append(kept, e)
|
||||
}
|
||||
open := openAt(t, kept, c.now())
|
||||
if len(open) != 4 {
|
||||
t.Fatalf("open %d", len(open))
|
||||
}
|
||||
h.Sync(open, c.now())
|
||||
replay(t, h, kept, nil)
|
||||
runFor(h, c, 30*time.Minute)
|
||||
// The urgent one, never said: one message on each channel urgent goes to. The warnings, hours
|
||||
// old: state, not news.
|
||||
if len(tg.sends) != 1 || len(dt.sends) != 1 || len(tg.edits)+len(dt.edits) != 0 {
|
||||
t.Fatalf("telegram %d desktop %d", len(tg.sends), len(dt.sends))
|
||||
}
|
||||
if !strings.HasPrefix(tg.sends[0].Title, "URGENT: ") || !strings.Contains(tg.sends[0].Text(), "merge.mesh-catalog.f2b51686.not-acted") {
|
||||
t.Fatalf("message: %q", tg.sends[0].Text())
|
||||
}
|
||||
if got := len(h.Open()); got != 4 {
|
||||
t.Fatalf("open here %d", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestAnEventOlderThanTheReadingIsNotActedOn(t *testing.T) {
|
||||
h, tg, dt, c, _ := newHolder(t)
|
||||
h.Sync(nil, c.now())
|
||||
k := cond("machine.ace.silent", Urgent, "the home server is silent")
|
||||
k.At = c.now().Add(-time.Minute)
|
||||
k.Raised = k.At
|
||||
h.Condition(EventRaised, k) // already cleared in the state read after it
|
||||
runFor(h, c, time.Minute)
|
||||
if said(tg, dt) != 0 || len(h.Open()) != 0 {
|
||||
t.Fatalf("a raising already over was said or kept")
|
||||
}
|
||||
}
|
||||
|
||||
func TestAnOldEventIsStateNeverAMessage(t *testing.T) {
|
||||
h, tg, dt, c, _ := newHolder(t)
|
||||
k := cond("machine.ace.silent", Urgent, "the home server is silent")
|
||||
k.At = c.now().Add(-FreshFor - time.Second)
|
||||
h.Condition(EventRaised, k)
|
||||
runFor(h, c, 2*time.Minute)
|
||||
if said(tg, dt) != 0 {
|
||||
t.Fatalf("an old raising was said")
|
||||
}
|
||||
open := h.Open()
|
||||
if len(open) != 1 || !open[0].History || !h.WantsSync() {
|
||||
t.Fatalf("not kept as state: %+v", open)
|
||||
}
|
||||
// Its clearing, now: it was never said, so nothing is.
|
||||
k.At = c.now()
|
||||
h.Condition(EventCleared, k)
|
||||
runFor(h, c, time.Minute)
|
||||
if said(tg, dt) != 0 {
|
||||
t.Fatalf("a clearing of what was never said was said")
|
||||
}
|
||||
// The state, read now, says an urgent one is open that was never said: that is said, once.
|
||||
h.Condition(EventRaised, Condition{Key: "machine.shanks.silent", Scope: "machine", ID: "shanks", Kind: "silent",
|
||||
Severity: Urgent, Summary: "the desktop is silent", At: c.now().Add(-time.Hour)})
|
||||
h.Sync([]Condition{cond("machine.shanks.silent", Urgent, "the desktop is silent")}, c.now())
|
||||
h.Sync([]Condition{cond("machine.shanks.silent", Urgent, "the desktop is silent")}, c.now().Add(time.Second))
|
||||
if len(tg.sends) != 1 || len(dt.sends) != 1 {
|
||||
t.Fatalf("telegram %d desktop %d", len(tg.sends), len(dt.sends))
|
||||
}
|
||||
}
|
||||
|
||||
func TestAClearingOfSomethingNeverSaidSaysNothing(t *testing.T) {
|
||||
h, tg, dt, c, _ := newHolder(t)
|
||||
k := cond("plan.41.stalled", Warning, "plan 41 waits")
|
||||
h.Condition(EventRaised, k)
|
||||
c.pass(10 * time.Second)
|
||||
h.Condition(EventCleared, k) // inside its burst window: it never went out
|
||||
runFor(h, c, 20*time.Minute)
|
||||
if said(tg, dt) != 0 {
|
||||
t.Fatalf("said: telegram %v desktop %v", tg.sends, dt.sends)
|
||||
}
|
||||
}
|
||||
|
||||
func TestABurstIsOneMessage(t *testing.T) {
|
||||
h, tg, _, c, _ := newHolder(t)
|
||||
h.Desktop = nil
|
||||
for i := 0; i < 5; i++ {
|
||||
h.Condition(EventRaised, cond(fmt.Sprintf("seat.s%d.silent", i), Warning, fmt.Sprintf("seat %d is silent", i)))
|
||||
c.pass(2 * time.Second)
|
||||
}
|
||||
h.Flush()
|
||||
if len(tg.sends) != 0 {
|
||||
t.Fatalf("sent inside the window: %d", len(tg.sends))
|
||||
}
|
||||
runFor(h, c, time.Minute)
|
||||
if len(tg.sends) != 1 || !strings.HasPrefix(tg.sends[0].Title, "WARNING: 5 new warnings") {
|
||||
t.Fatalf("sends %d: %+v", len(tg.sends), tg.sends)
|
||||
}
|
||||
for i := 0; i < 5; i++ {
|
||||
if !strings.Contains(tg.sends[0].Body, fmt.Sprintf("seat.s%d.silent", i)) {
|
||||
t.Fatalf("the digest does not name s%d: %q", i, tg.sends[0].Body)
|
||||
}
|
||||
}
|
||||
// Said in a digest, a clearing is a line in the next one — never an edit of the digest.
|
||||
for i := 0; i < 3; i++ {
|
||||
h.Condition(EventCleared, Condition{Key: fmt.Sprintf("seat.s%d.silent", i)})
|
||||
}
|
||||
runFor(h, c, time.Minute)
|
||||
if len(tg.sends) != 2 || len(tg.edits) != 0 || !strings.HasPrefix(tg.sends[1].Title, "CLEARED: 3 cleared") {
|
||||
t.Fatalf("sends %d edits %d: %q", len(tg.sends), len(tg.edits), tg.sends[len(tg.sends)-1].Title)
|
||||
}
|
||||
// A burst of urgent ones: the first at once, the rest as one.
|
||||
c.pass(time.Minute)
|
||||
for i := 0; i < 6; i++ {
|
||||
h.Condition(EventRaised, cond(fmt.Sprintf("machine.m%d.silent", i), Urgent, "a machine is silent"))
|
||||
c.pass(time.Second)
|
||||
}
|
||||
runFor(h, c, time.Minute)
|
||||
if len(tg.sends) != 4 || !strings.HasPrefix(tg.sends[3].Title, "URGENT: 5 new urgent") {
|
||||
t.Fatalf("sends %d: %q", len(tg.sends), tg.sends[len(tg.sends)-1].Title)
|
||||
}
|
||||
}
|
||||
|
||||
func TestTheDesktopIsGentle(t *testing.T) {
|
||||
h, tg, dt, c, _ := newHolder(t)
|
||||
h.Condition(EventRaised, cond("plan.1.stalled", Warning, "plan 1 waits"))
|
||||
h.Condition(EventRaised, cond("plan.2.stalled", Warning, "plan 2 waits"))
|
||||
runFor(h, c, time.Minute)
|
||||
if len(dt.sends) != 1 || len(tg.sends) != 0 {
|
||||
t.Fatalf("desktop %d telegram %d", len(dt.sends), len(tg.sends))
|
||||
}
|
||||
h.Condition(EventRaised, cond("plan.3.stalled", Warning, "plan 3 waits"))
|
||||
runFor(h, c, 5*time.Minute)
|
||||
h.Condition(EventRaised, cond("plan.4.stalled", Warning, "plan 4 waits"))
|
||||
// Urgent is not held to the desktop's rhythm.
|
||||
h.Condition(EventRaised, cond("machine.ace.silent", Urgent, "the home server is silent"))
|
||||
h.Flush()
|
||||
if len(dt.sends) != 2 || !strings.HasPrefix(dt.sends[1].Title, "URGENT: ") || len(tg.sends) != 1 {
|
||||
t.Fatalf("urgent: desktop %d telegram %d", len(dt.sends), len(tg.sends))
|
||||
}
|
||||
// The urgent message took the held warnings along — one notification, read in one go. What
|
||||
// comes after waits for the fifteen minutes.
|
||||
h.Condition(EventRaised, cond("plan.5.stalled", Warning, "plan 5 waits"))
|
||||
runFor(h, c, 8*time.Minute)
|
||||
if len(dt.sends) != 2 {
|
||||
t.Fatalf("a warning reached the desktop inside fifteen minutes: %d", len(dt.sends))
|
||||
}
|
||||
runFor(h, c, 15*time.Minute)
|
||||
if len(dt.sends) != 3 {
|
||||
t.Fatalf("desktop %d", len(dt.sends))
|
||||
}
|
||||
for i := 0; i < len(dt.sends)-1; i++ {
|
||||
if !dt.sends[i].Urgent && !dt.sends[i+1].Urgent {
|
||||
t.Fatalf("two warning messages close together")
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestNothingIsSentInTheFirstMinuteButTheUrgentOpenAsOne(t *testing.T) {
|
||||
h, tg, dt, c, _ := newHolder(t)
|
||||
h.Start()
|
||||
u1 := cond("machine.ace.silent", Urgent, "the home server is silent")
|
||||
u1.Raised = c.now().Add(-3 * time.Hour)
|
||||
u2 := cond("bus.controller.refused", Urgent, "the bus refuses the controller")
|
||||
u2.Raised = c.now().Add(-20 * time.Minute)
|
||||
old := cond("plan.41.stalled", Warning, "plan 41 waits")
|
||||
old.Raised = c.now().Add(-2 * time.Hour)
|
||||
fresh := cond("plan.42.stalled", Warning, "plan 42 waits")
|
||||
fresh.Raised = c.now().Add(-2 * time.Minute)
|
||||
res := h.Sync([]Condition{u1, u2, old, fresh}, c.now())
|
||||
if len(res.Said) != 2 || len(res.State) != 1 || len(res.News) != 1 {
|
||||
t.Fatalf("sync: %+v", res)
|
||||
}
|
||||
if len(tg.sends) != 1 || len(dt.sends) != 1 || !strings.HasPrefix(tg.sends[0].Title, "URGENT: 2 new urgent") {
|
||||
t.Fatalf("telegram %d desktop %d: %+v", len(tg.sends), len(dt.sends), tg.sends)
|
||||
}
|
||||
c.pass(10 * time.Second)
|
||||
k := cond("plan.43.stalled", Warning, "plan 43 waits")
|
||||
k.At = c.now()
|
||||
h.Condition(EventRaised, k)
|
||||
u3 := cond("machine.shanks.silent", Urgent, "the desktop is silent")
|
||||
u3.At = c.now()
|
||||
h.Condition(EventRaised, u3)
|
||||
runFor(h, c, 45*time.Second)
|
||||
if len(tg.sends) != 1 || len(dt.sends) != 1 {
|
||||
t.Fatalf("sent in the grace minute: telegram %d desktop %d", len(tg.sends), len(dt.sends))
|
||||
}
|
||||
runFor(h, c, 20*time.Second)
|
||||
// After it: what waited, as one message per channel.
|
||||
if len(dt.sends) != 2 || len(tg.sends) != 2 {
|
||||
t.Fatalf("after the grace: telegram %d desktop %d", len(tg.sends), len(dt.sends))
|
||||
}
|
||||
if d := dt.sends[1]; !strings.Contains(d.Body, "plan.42.stalled") || !strings.Contains(d.Body, "plan.43.stalled") ||
|
||||
!strings.Contains(d.Body, "machine.shanks.silent") || strings.Contains(d.Body, "plan.41.stalled") {
|
||||
t.Fatalf("desktop digest: %q", d.Text())
|
||||
}
|
||||
}
|
||||
|
||||
func TestWhatEndedWhileAwayIsEndedQuietly(t *testing.T) {
|
||||
h, tg, dt, c, _ := newHolder(t)
|
||||
k := cond("machine.ace.silent", Urgent, "the home server is silent")
|
||||
k.At = c.now()
|
||||
h.Condition(EventRaised, k)
|
||||
if len(tg.sends) != 1 || len(dt.sends) != 1 {
|
||||
t.Fatalf("telegram %d desktop %d", len(tg.sends), len(dt.sends))
|
||||
}
|
||||
c.pass(2 * time.Hour)
|
||||
h.Sync(nil, c.now())
|
||||
runFor(h, c, 20*time.Minute)
|
||||
// Telegram's edit notifies nobody, so the message there says it is over; the desktop is not
|
||||
// woken for something that ended while nobody said so.
|
||||
if len(tg.sends) != 1 || len(dt.sends) != 1 || len(tg.edits) != 1 || len(dt.edits) != 0 {
|
||||
t.Fatalf("telegram %d/%d desktop %d/%d", len(tg.sends), len(tg.edits), len(dt.sends), len(dt.edits))
|
||||
}
|
||||
if len(h.Open()) != 0 {
|
||||
t.Fatalf("still open here")
|
||||
}
|
||||
}
|
||||
|
||||
func TestARestartSaysNothingAgain(t *testing.T) {
|
||||
h, tg, dt, c, st := newHolder(t)
|
||||
k := cond("machine.ace.silent", Urgent, "the home server is silent")
|
||||
k.At = c.now()
|
||||
h.Condition(EventRaised, k)
|
||||
w := cond("plan.41.stalled", Warning, "plan 41 waits")
|
||||
w.At = c.now()
|
||||
h.Condition(EventRaised, w) // held when the holder stops
|
||||
c.pass(5 * time.Second)
|
||||
again := &Holder{Telegram: tg, Desktop: dt, Store: st, Now: c.now, Logf: t.Logf}
|
||||
if err := again.Load(); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
again.Start()
|
||||
again.Sync([]Condition{k, w}, c.now())
|
||||
runFor(again, c, time.Hour)
|
||||
if len(tg.sends) != 1 || len(dt.sends) != 1 {
|
||||
t.Fatalf("a restart said it again: telegram %d desktop %d", len(tg.sends), len(dt.sends))
|
||||
}
|
||||
}
|
||||
|
||||
func TestTheOpenSetIsReadInEveryShape(t *testing.T) {
|
||||
one := `{"key":"machine.ace.silent","kind":"silent","subject":{"scope":"machine","id":"ace"},"severity":"urgent","summary":"s","raised":"2026-10-06T09:00:00Z"}`
|
||||
for _, raw := range []string{
|
||||
`{"conditions":[` + one + `],"open":1}`,
|
||||
`{"answer":{"conditions":[` + one + `]},"ok":true}`,
|
||||
`{"output":"{\"conditions\":[` + strings.ReplaceAll(one, `"`, `\"`) + `]}"}`,
|
||||
`{"content":[{"type":"text","text":"{\"conditions\":[` + strings.ReplaceAll(one, `"`, `\"`) + `]}"}]}`,
|
||||
} {
|
||||
got, err := DecodeOpen(json.RawMessage(raw))
|
||||
if err != nil || len(got) != 1 || got[0].Key != "machine.ace.silent" || got[0].Severity != Urgent {
|
||||
t.Fatalf("%s: %+v %v", raw, got, err)
|
||||
}
|
||||
}
|
||||
if got, err := DecodeOpen(json.RawMessage(`{"conditions":[],"open":0}`)); err != nil || len(got) != 0 {
|
||||
t.Fatalf("empty: %v %v", got, err)
|
||||
}
|
||||
for _, bad := range []string{`{"ok":true}`, `[]`, `{"content":[{"text":"denied"}],"isError":true}`} {
|
||||
if _, err := DecodeOpen(json.RawMessage(bad)); err == nil {
|
||||
t.Fatalf("%s was read as an open set", bad)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestAnEventSaysItsOwnTime(t *testing.T) {
|
||||
c, err := DecodeCondition(EventCleared, []byte(`{"key":"machine.ace.silent","at":"2026-10-06T08:36:50Z","raised":"2026-10-06T08:32:30Z"}`))
|
||||
if err != nil || !c.When(EventCleared).Equal(time.Date(2026, 10, 6, 8, 36, 50, 0, time.UTC)) {
|
||||
t.Fatalf("%v %v", c.When(EventCleared), err)
|
||||
}
|
||||
c.At = time.Time{}
|
||||
if !c.When(EventRaised).Equal(c.Raised) || !c.When(EventChanged).IsZero() {
|
||||
t.Fatalf("fallbacks")
|
||||
}
|
||||
}
|
||||
@@ -3,21 +3,47 @@ package main
|
||||
// The holder of the operator-channel seat (novox/hq to-be 45 §5, ADR 0227): what is sent, to whom,
|
||||
// when, and how often. The controller decides what is wrong; this decides what is said.
|
||||
//
|
||||
// **History is never news** (novox/hq issue 271). The events stream keeps what it was given and hands
|
||||
// all of it to a consumer that is new, recreated or was away; on 2026-10-06 a newly made consumer
|
||||
// replayed three hours of raised-and-cleared conditions and each was said as if it had just happened.
|
||||
// So what is said is decided by the state now and by when a thing happened, never by when it arrived:
|
||||
//
|
||||
// - On start, and whenever it may have missed something (every ten minutes, and at once when an
|
||||
// old event arrives), the holder reads what is open now from the controller (`conditions`) and
|
||||
// takes it as state: nothing it says is sent for what it already said, what cleared meanwhile is
|
||||
// forgotten without a word (a Telegram message is edited, which notifies nobody), an urgent
|
||||
// condition open and never said is said — all of them in one message — and a warning raised in
|
||||
// the last ten minutes is said like any other; an older warning is state.
|
||||
// - An event older than that reading is already in it, and is not acted on. An event older than
|
||||
// FreshFor (ten minutes) by its own time is state only: it is recorded, never said.
|
||||
// - A clearing of something never said says nothing; a held message whose condition cleared before
|
||||
// it went out is dropped.
|
||||
// - Nothing is sent in the first minute after start (StartGrace) but the one message of urgent
|
||||
// conditions open in the controller's state; what arrives meanwhile waits for the minute's end.
|
||||
//
|
||||
// What is said, and how often:
|
||||
//
|
||||
// - On `condition-raised`: one message, deduplicated by the condition's key. Said again with the
|
||||
// same key while open, it is the same message, not a second.
|
||||
// - Once more if still open after 1 hour (urgent) or 12 hours (warning).
|
||||
// - On `condition-cleared`: the first message is edited where the channel can (both can);
|
||||
// otherwise a new one says it. Cleared and raised again within ten minutes, it is the same
|
||||
// message, edited back to open — not a new one.
|
||||
// - On `condition-cleared`: the first message is edited where the channel can (both can).
|
||||
// Cleared and raised again within ten minutes, it is the same message, edited back to open.
|
||||
// - A silenced condition sends nothing.
|
||||
// - Urgent to both channels; warning to the desktop when the operator's session is there,
|
||||
// otherwise to Telegram.
|
||||
// - At most twenty messages an hour per channel; the excess is held and folded into one message
|
||||
// naming them all, sent at most every ten minutes — the cap is said, never silent.
|
||||
// - **Bursts are one message.** Messages for a channel are held for BurstWindow (30 s) from the
|
||||
// first and go out together: one alone is sent as itself, several as one digest ("5 new
|
||||
// warnings: …"). An urgent message is sent at once when nothing went out on its channel in the
|
||||
// last 30 s.
|
||||
// - **The desktop is gentle:** warnings reach it as at most one message every DesktopWarningEvery
|
||||
// (15 min); urgent ones are not held to that.
|
||||
// - At most twenty messages an hour per channel — a last line of defence, since the above keeps a
|
||||
// channel far below it. Reached, it is said once in a message of its own and the rest is held,
|
||||
// to go out as one digest when the hour allows.
|
||||
// - A message carrying an address, a path or a secret is refused (content.go); what is sent in its
|
||||
// place says `channel-refused`, with the offending part withheld.
|
||||
// - A channel that cannot send says so in the status and the log, and the message is tried again
|
||||
// every minute while the condition is open.
|
||||
// - A channel that cannot send says so in the status and the log, and is tried again every
|
||||
// minute while what it holds is still open.
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
@@ -32,9 +58,26 @@ const (
|
||||
RemindWarning = 12 * time.Hour
|
||||
ReopenWindow = 10 * time.Minute
|
||||
CapPerHour = 20
|
||||
FoldEvery = 10 * time.Minute
|
||||
KeptSends = 200
|
||||
KeptRefusals = 50
|
||||
|
||||
// FreshFor: an event older than this by its own time is history — state, never a message.
|
||||
FreshFor = 10 * time.Minute
|
||||
// BurstWindow: what arrives for a channel within it goes out as one message.
|
||||
BurstWindow = 30 * time.Second
|
||||
// DesktopWarningEvery: warnings reach the desktop at most this often.
|
||||
DesktopWarningEvery = 15 * time.Minute
|
||||
// StartGrace: after start, nothing but the urgent conditions open in the controller's state.
|
||||
StartGrace = time.Minute
|
||||
// ResyncEvery: how often the open set is read from the controller again.
|
||||
ResyncEvery = 10 * time.Minute
|
||||
// RetryEvery: a channel that failed is tried again after this.
|
||||
RetryEvery = time.Minute
|
||||
// ClockSlack: the controller's clock and this one may differ by this much; an event this close
|
||||
// to the reading is taken rather than assumed to be in it (taking one twice is harmless).
|
||||
ClockSlack = 5 * time.Second
|
||||
// DigestLines: a digest names at most this many, then says how many more.
|
||||
DigestLines = 15
|
||||
)
|
||||
|
||||
// Message is what a channel shows: a title line and a body.
|
||||
@@ -61,6 +104,15 @@ type Channel interface {
|
||||
CanEdit() bool
|
||||
}
|
||||
|
||||
// silentEditor is a channel whose edits notify nobody (Telegram's are): what is over is said there
|
||||
// even when it ended while the holder was away, since saying it costs the operator nothing.
|
||||
type silentEditor interface{ SilentEdit() bool }
|
||||
|
||||
func editsSilently(ch Channel) bool {
|
||||
s, ok := ch.(silentEditor)
|
||||
return ok && s.SilentEdit()
|
||||
}
|
||||
|
||||
// Record is one open message, kept in the module's own state so a restart forgets nothing.
|
||||
type Record struct {
|
||||
Key string `json:"key"`
|
||||
@@ -72,27 +124,33 @@ type Record struct {
|
||||
More string `json:"more"`
|
||||
Raised time.Time `json:"raised"`
|
||||
SilencedTill time.Time `json:"silenced_till,omitempty"`
|
||||
Sent map[string]string `json:"sent,omitempty"` // channel -> the first message's id
|
||||
Sent map[string]string `json:"sent,omitempty"` // channel -> the first message's id; "" when said in a digest
|
||||
FirstSent time.Time `json:"first_sent,omitempty"`
|
||||
Reminded bool `json:"reminded,omitempty"`
|
||||
Pending string `json:"pending,omitempty"` // what is still to be said: raised, reminder, …
|
||||
Folded bool `json:"folded,omitempty"`
|
||||
Pending string `json:"pending,omitempty"` // what could not be said yet: raised, reminder, …
|
||||
Folded bool `json:"folded,omitempty"` // held by the cap
|
||||
Cleared time.Time `json:"cleared,omitempty"`
|
||||
Count int `json:"count"`
|
||||
Refused string `json:"refused,omitempty"`
|
||||
// History: known from the state or from an old event, and never said — not news.
|
||||
History bool `json:"history,omitempty"`
|
||||
// Seen is when this holder last took word of it; the open set read after it is authoritative.
|
||||
Seen time.Time `json:"seen,omitempty"`
|
||||
}
|
||||
|
||||
func (r *Record) silenced(now time.Time) bool {
|
||||
return !r.SilencedTill.IsZero() && now.Before(r.SilencedTill)
|
||||
}
|
||||
|
||||
func (r *Record) told() bool { return len(r.Sent) > 0 }
|
||||
|
||||
// Sent is one message that went out, or was held, for the history.
|
||||
type Sent struct {
|
||||
At time.Time `json:"at"`
|
||||
Channel string `json:"channel"`
|
||||
Key string `json:"key"`
|
||||
What string `json:"what"`
|
||||
Outcome string `json:"outcome"` // sent, edited, folded, failed: <why>
|
||||
Outcome string `json:"outcome"` // sent, edited, held, in digest, failed: <why>
|
||||
}
|
||||
|
||||
// RefusalNote is one refused message: never its text.
|
||||
@@ -112,10 +170,6 @@ type Store interface {
|
||||
Recent() ([]Sent, error)
|
||||
}
|
||||
|
||||
type foldEntry struct {
|
||||
Key, Severity, What string
|
||||
}
|
||||
|
||||
// Holder is the seat's holder.
|
||||
type Holder struct {
|
||||
Telegram Channel
|
||||
@@ -132,8 +186,11 @@ type Holder struct {
|
||||
open map[string]*Record
|
||||
recent []Sent
|
||||
refusals []RefusalNote
|
||||
folds map[string][]foldEntry
|
||||
lastFold map[string]time.Time
|
||||
outbox map[string][]item
|
||||
lastSend map[string]time.Time
|
||||
lastWarn map[string]time.Time
|
||||
lastTry map[string]time.Time
|
||||
capSaid map[string]time.Time
|
||||
chanErr map[string]string
|
||||
chanErrAt map[string]time.Time
|
||||
chanOK map[string]time.Time
|
||||
@@ -144,13 +201,22 @@ type Holder struct {
|
||||
storeErr string
|
||||
heard map[string]int
|
||||
lastHeard time.Time
|
||||
started time.Time
|
||||
syncedAt time.Time
|
||||
syncErr string
|
||||
wantSync bool
|
||||
history int
|
||||
stale int
|
||||
}
|
||||
|
||||
func (h *Holder) init() {
|
||||
if h.open == nil {
|
||||
h.open = map[string]*Record{}
|
||||
h.folds = map[string][]foldEntry{}
|
||||
h.lastFold = map[string]time.Time{}
|
||||
h.outbox = map[string][]item{}
|
||||
h.lastSend = map[string]time.Time{}
|
||||
h.lastWarn = map[string]time.Time{}
|
||||
h.lastTry = map[string]time.Time{}
|
||||
h.capSaid = map[string]time.Time{}
|
||||
h.chanErr = map[string]string{}
|
||||
h.chanErrAt = map[string]time.Time{}
|
||||
h.chanOK = map[string]time.Time{}
|
||||
@@ -165,7 +231,8 @@ func (h *Holder) init() {
|
||||
}
|
||||
}
|
||||
|
||||
// Load reads back what was open and what was sent before a restart.
|
||||
// Load reads back what was open and what was sent before a restart. What was waiting to be said is
|
||||
// not said now: it is history, and the controller's state, read next, decides what is still news.
|
||||
func (h *Holder) Load() error {
|
||||
h.work.Lock()
|
||||
defer h.work.Unlock()
|
||||
@@ -182,6 +249,12 @@ func (h *Holder) Load() error {
|
||||
}
|
||||
for i := range recs {
|
||||
r := recs[i]
|
||||
if r.Pending != "" || r.Folded {
|
||||
r.Pending, r.Folded = "", false
|
||||
if !r.told() {
|
||||
r.History = true
|
||||
}
|
||||
}
|
||||
h.open[r.Key] = &r
|
||||
}
|
||||
if sent, err := h.Store.Recent(); err == nil {
|
||||
@@ -190,26 +263,66 @@ func (h *Holder) Load() error {
|
||||
return nil
|
||||
}
|
||||
|
||||
// Condition takes one of the controller's condition events.
|
||||
// Start marks the holder's start: the grace minute counts from here.
|
||||
func (h *Holder) Start() {
|
||||
h.work.Lock()
|
||||
defer h.work.Unlock()
|
||||
h.mu.Lock()
|
||||
defer h.mu.Unlock()
|
||||
h.init()
|
||||
h.started = h.Now()
|
||||
}
|
||||
|
||||
func (h *Holder) inGrace(now time.Time) bool {
|
||||
return !h.started.IsZero() && now.Sub(h.started) < StartGrace
|
||||
}
|
||||
|
||||
// Condition takes one of the controller's condition events. Whether it is news is decided by its own
|
||||
// time, never by its arrival.
|
||||
func (h *Holder) Condition(event string, c Condition) {
|
||||
h.work.Lock()
|
||||
defer h.work.Unlock()
|
||||
h.mu.Lock()
|
||||
h.init()
|
||||
now := h.Now()
|
||||
h.heard[event]++
|
||||
h.lastHeard = h.Now()
|
||||
h.lastHeard = now
|
||||
synced := h.syncedAt
|
||||
h.mu.Unlock()
|
||||
rec := Record{
|
||||
Key: c.Key, Kind: c.Kind, Subject: c.SubjectWords(), Severity: c.Severity, Summary: c.Summary,
|
||||
Origin: "condition", More: "conditions show " + c.Key, Raised: c.Raised, SilencedTill: c.SilencedTill,
|
||||
}
|
||||
when := c.When(event)
|
||||
quiet := false
|
||||
switch {
|
||||
case !when.IsZero() && !synced.IsZero() && when.Before(synced.Add(-ClockSlack)):
|
||||
// Already in the state read from the controller after it: nothing to take.
|
||||
h.mu.Lock()
|
||||
h.history++
|
||||
h.mu.Unlock()
|
||||
return
|
||||
case !when.IsZero() && now.Sub(when) > FreshFor:
|
||||
// Old: a replay, a redelivery, or the holder was away. State, never a message — and the
|
||||
// controller's state is read again, since what came between may be missing too.
|
||||
quiet = true
|
||||
h.mu.Lock()
|
||||
h.stale++
|
||||
first := !h.wantSync
|
||||
h.wantSync = true
|
||||
h.mu.Unlock()
|
||||
if first {
|
||||
h.Logf("[messenger] %s %s happened %s ago: taken as state, not news; reading the open conditions again",
|
||||
event, c.Key, roughly(now.Sub(when)))
|
||||
}
|
||||
}
|
||||
switch event {
|
||||
case EventRaised:
|
||||
h.raised(rec)
|
||||
h.raised(rec, quiet)
|
||||
case EventChanged:
|
||||
h.changed(rec)
|
||||
h.changed(rec, quiet)
|
||||
case EventCleared:
|
||||
h.cleared(rec.Key)
|
||||
h.cleared(rec.Key, quiet)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -232,7 +345,7 @@ func (h *Holder) Unreadable(key string, err error) {
|
||||
})
|
||||
}
|
||||
|
||||
func (h *Holder) raised(rec Record) {
|
||||
func (h *Holder) raised(rec Record, quiet bool) {
|
||||
now := h.Now()
|
||||
h.mu.Lock()
|
||||
old := h.open[rec.Key]
|
||||
@@ -244,6 +357,7 @@ func (h *Holder) raised(rec Record) {
|
||||
if rec.Severity != "" {
|
||||
old.Severity = rec.Severity
|
||||
}
|
||||
old.Seen = now
|
||||
h.mu.Unlock()
|
||||
h.persist(old)
|
||||
return
|
||||
@@ -254,9 +368,10 @@ func (h *Holder) raised(rec Record) {
|
||||
old.Cleared = time.Time{}
|
||||
old.Count++
|
||||
old.Summary, old.Severity, old.SilencedTill = rec.Summary, rec.Severity, rec.SilencedTill
|
||||
old.Seen = now
|
||||
h.mu.Unlock()
|
||||
if !old.silenced(now) && len(old.Sent) > 0 {
|
||||
h.edit(old, "reopened")
|
||||
if !quiet && !old.silenced(now) && old.told() {
|
||||
h.edit(old, "reopened", false)
|
||||
}
|
||||
h.persist(old)
|
||||
return
|
||||
@@ -267,26 +382,30 @@ func (h *Holder) raised(rec Record) {
|
||||
r.Raised = now
|
||||
}
|
||||
r.Sent = map[string]string{}
|
||||
r.Seen = now
|
||||
r.History = quiet
|
||||
h.mu.Lock()
|
||||
h.open[r.Key] = &r
|
||||
h.mu.Unlock()
|
||||
if r.silenced(now) {
|
||||
switch {
|
||||
case quiet:
|
||||
case r.silenced(now):
|
||||
h.Logf("[messenger] %s raised while silenced until %s: nothing sent", r.Key, r.SilencedTill.Format(time.RFC3339))
|
||||
h.persist(&r)
|
||||
return
|
||||
default:
|
||||
h.deliver(&r, "raised")
|
||||
}
|
||||
h.deliver(&r, "raised")
|
||||
h.persist(&r)
|
||||
}
|
||||
|
||||
func (h *Holder) changed(rec Record) {
|
||||
func (h *Holder) changed(rec Record, quiet bool) {
|
||||
h.mu.Lock()
|
||||
old := h.open[rec.Key]
|
||||
h.mu.Unlock()
|
||||
if old == nil || !old.Cleared.IsZero() {
|
||||
// A change to something this holder never heard raised: read as a raise, so it is said.
|
||||
// A change to something this holder never heard raised: read as a raise, so it is said
|
||||
// (when it is news).
|
||||
h.Logf("[messenger] %s changed and was not open here; taken as raised", rec.Key)
|
||||
h.raised(rec)
|
||||
h.raised(rec, quiet)
|
||||
return
|
||||
}
|
||||
h.mu.Lock()
|
||||
@@ -295,15 +414,19 @@ func (h *Holder) changed(rec Record) {
|
||||
if rec.Severity != "" {
|
||||
old.Severity = rec.Severity
|
||||
}
|
||||
old.Seen = h.Now()
|
||||
h.mu.Unlock()
|
||||
if escalated && !old.silenced(h.Now()) {
|
||||
if escalated && !quiet && !old.silenced(h.Now()) {
|
||||
// Routing differs for urgent: said once more, to both channels.
|
||||
h.deliver(old, "escalated")
|
||||
}
|
||||
h.persist(old)
|
||||
}
|
||||
|
||||
func (h *Holder) cleared(key string) {
|
||||
// cleared ends a record. Something never said is ended without a word; what was waiting to be said
|
||||
// about it is dropped. quiet: it ended while this holder was not listening, so it is said only where
|
||||
// an edit notifies nobody.
|
||||
func (h *Holder) cleared(key string, quiet bool) {
|
||||
now := h.Now()
|
||||
h.mu.Lock()
|
||||
old := h.open[key]
|
||||
@@ -314,27 +437,27 @@ func (h *Holder) cleared(key string) {
|
||||
}
|
||||
h.mu.Lock()
|
||||
old.Cleared = now
|
||||
pending := old.Pending
|
||||
old.Pending = ""
|
||||
sent := len(old.Sent) > 0
|
||||
folded := old.Folded
|
||||
old.Folded = false
|
||||
old.Seen = now
|
||||
told := old.told()
|
||||
h.mu.Unlock()
|
||||
dropped := h.unqueue(key)
|
||||
switch {
|
||||
case old.silenced(now):
|
||||
// A silenced condition sends nothing, its clearing included.
|
||||
case sent:
|
||||
h.edit(old, "cleared")
|
||||
case folded:
|
||||
// Held by the cap and never sent on its own: its clearing is said like any message.
|
||||
h.deliver(old, "cleared")
|
||||
case pending != "":
|
||||
h.Logf("[messenger] %s cleared before it could be sent: nothing to unsay", key)
|
||||
case !told:
|
||||
if dropped > 0 {
|
||||
h.Logf("[messenger] %s cleared before it was said: nothing to unsay", key)
|
||||
}
|
||||
default:
|
||||
h.edit(old, "cleared", quiet)
|
||||
}
|
||||
h.persist(old)
|
||||
}
|
||||
|
||||
// Tick does what time asks: reminders, retries, the folded message, and forgetting what cleared
|
||||
// long enough ago that a new raise is a new message.
|
||||
// Tick does what time asks: reminders, retries, what is held, and forgetting what cleared long
|
||||
// enough ago that a new raise is a new message.
|
||||
func (h *Holder) Tick() {
|
||||
h.work.Lock()
|
||||
defer h.work.Unlock()
|
||||
@@ -361,10 +484,14 @@ func (h *Holder) Tick() {
|
||||
}
|
||||
}
|
||||
case r.silenced(now):
|
||||
case r.Pending != "":
|
||||
case r.Pending != "" && !h.queued(r.Key):
|
||||
// No channel could take it when it was said: tried again.
|
||||
h.deliver(r, r.Pending)
|
||||
h.persist(r)
|
||||
case !r.Reminded && !r.FirstSent.IsZero() && now.Sub(r.Raised) >= remindAfter(r.Severity):
|
||||
case !r.Reminded && !r.FirstSent.IsZero() && now.Sub(r.Raised) >= remindAfter(r.Severity) &&
|
||||
now.Sub(r.FirstSent) >= remindAfter(r.Severity):
|
||||
// Once more after its bound — counted from its raising, and never sooner after it was
|
||||
// first said (an old condition said at a reading of the state is not reminded at once).
|
||||
h.mu.Lock()
|
||||
r.Reminded = true
|
||||
h.mu.Unlock()
|
||||
@@ -372,7 +499,19 @@ func (h *Holder) Tick() {
|
||||
h.persist(r)
|
||||
}
|
||||
}
|
||||
h.flushFolds(now)
|
||||
h.flush(now, false)
|
||||
}
|
||||
|
||||
// Flush sends what is due on each channel; the runtime calls it every few seconds, so a held
|
||||
// message waits its window and not a minute more.
|
||||
func (h *Holder) Flush() {
|
||||
h.work.Lock()
|
||||
defer h.work.Unlock()
|
||||
h.mu.Lock()
|
||||
h.init()
|
||||
now := h.Now()
|
||||
h.mu.Unlock()
|
||||
h.flush(now, false)
|
||||
}
|
||||
|
||||
func remindAfter(severity string) time.Duration {
|
||||
@@ -476,39 +615,6 @@ func (h *Holder) say(r *Record, what string) Message {
|
||||
}
|
||||
|
||||
// deliver sends a record's message where its severity routes it.
|
||||
func (h *Holder) deliver(r *Record, what string) {
|
||||
m := h.say(r, what)
|
||||
now := h.Now()
|
||||
delivered := false
|
||||
if r.Severity == Urgent {
|
||||
for _, ch := range h.channels() {
|
||||
if h.sendOn(ch, r, what, m) {
|
||||
delivered = true
|
||||
}
|
||||
}
|
||||
} else {
|
||||
if h.Desktop != nil && h.Desktop.Ready() == nil {
|
||||
delivered = h.sendOn(h.Desktop, r, what, m)
|
||||
}
|
||||
if !delivered && h.Telegram != nil {
|
||||
delivered = h.sendOn(h.Telegram, r, what, m)
|
||||
}
|
||||
}
|
||||
h.mu.Lock()
|
||||
if delivered {
|
||||
r.Pending = ""
|
||||
if r.FirstSent.IsZero() && (what == "raised" || what == "escalated") {
|
||||
r.FirstSent = now
|
||||
}
|
||||
} else {
|
||||
// Nothing took it: said in the status and the log, tried again next minute.
|
||||
r.Pending = what
|
||||
}
|
||||
h.mu.Unlock()
|
||||
if !delivered {
|
||||
h.Logf("[messenger] could not send %s %s on any channel; trying again every minute: %s", what, r.Key, h.whyNot())
|
||||
}
|
||||
}
|
||||
|
||||
func (h *Holder) channels() []Channel {
|
||||
var out []Channel
|
||||
@@ -520,169 +626,6 @@ func (h *Holder) channels() []Channel {
|
||||
return out
|
||||
}
|
||||
|
||||
// sendOn sends one message on one channel under its cap; held by the cap, it is folded, which counts
|
||||
// as delivered — the fold will say it.
|
||||
func (h *Holder) sendOn(ch Channel, r *Record, what string, m Message) bool {
|
||||
if err := ch.Ready(); err != nil {
|
||||
h.noteChannel(ch.Name(), err)
|
||||
return false
|
||||
}
|
||||
if !h.allow(ch.Name()) {
|
||||
h.foldOn(ch.Name(), r, what)
|
||||
return true
|
||||
}
|
||||
id, err := ch.Send(m)
|
||||
h.noteChannel(ch.Name(), err)
|
||||
if err != nil {
|
||||
h.record(Sent{At: h.Now(), Channel: ch.Name(), Key: r.Key, What: what, Outcome: "failed: " + err.Error()})
|
||||
return false
|
||||
}
|
||||
h.mu.Lock()
|
||||
if r.Sent == nil {
|
||||
r.Sent = map[string]string{}
|
||||
}
|
||||
if _, has := r.Sent[ch.Name()]; !has && what != "cleared" {
|
||||
r.Sent[ch.Name()] = id
|
||||
}
|
||||
h.mu.Unlock()
|
||||
h.record(Sent{At: h.Now(), Channel: ch.Name(), Key: r.Key, What: what, Outcome: "sent"})
|
||||
return true
|
||||
}
|
||||
|
||||
// edit changes the first message on each channel that showed it; a channel that cannot, or whose
|
||||
// edit fails, is sent a new message instead.
|
||||
func (h *Holder) edit(r *Record, what string) {
|
||||
m := h.say(r, what)
|
||||
h.mu.Lock()
|
||||
sent := map[string]string{}
|
||||
for k, v := range r.Sent {
|
||||
sent[k] = v
|
||||
}
|
||||
h.mu.Unlock()
|
||||
for _, ch := range h.channels() {
|
||||
id, shown := sent[ch.Name()]
|
||||
if !shown {
|
||||
continue
|
||||
}
|
||||
if ch.CanEdit() {
|
||||
err := ch.Edit(id, m)
|
||||
h.noteChannel(ch.Name(), err)
|
||||
if err == nil {
|
||||
h.record(Sent{At: h.Now(), Channel: ch.Name(), Key: r.Key, What: what, Outcome: "edited"})
|
||||
continue
|
||||
}
|
||||
h.Logf("[messenger] could not edit the message for %s on %s (%v); sending a new one", r.Key, ch.Name(), err)
|
||||
}
|
||||
h.sendOn(ch, r, what, m)
|
||||
}
|
||||
}
|
||||
|
||||
func (h *Holder) foldOn(channel string, r *Record, what string) {
|
||||
h.mu.Lock()
|
||||
r.Folded = true
|
||||
list := h.folds[channel]
|
||||
replaced := false
|
||||
for i := range list {
|
||||
if list[i].Key == r.Key {
|
||||
list[i].What, list[i].Severity, replaced = what, r.Severity, true
|
||||
}
|
||||
}
|
||||
if !replaced {
|
||||
list = append(list, foldEntry{Key: r.Key, Severity: r.Severity, What: what})
|
||||
}
|
||||
h.folds[channel] = list
|
||||
first := len(list) == 1 && !replaced
|
||||
h.mu.Unlock()
|
||||
if first {
|
||||
h.Logf("[messenger] %s is at its cap of %d messages an hour: holding the rest, to be folded into one", channel, CapPerHour)
|
||||
}
|
||||
h.record(Sent{At: h.Now(), Channel: channel, Key: r.Key, What: what, Outcome: "folded"})
|
||||
}
|
||||
|
||||
func (h *Holder) flushFolds(now time.Time) {
|
||||
for _, ch := range h.channels() {
|
||||
h.mu.Lock()
|
||||
list := append([]foldEntry(nil), h.folds[ch.Name()]...)
|
||||
last := h.lastFold[ch.Name()]
|
||||
h.mu.Unlock()
|
||||
if len(list) == 0 || now.Sub(last) < FoldEvery {
|
||||
continue
|
||||
}
|
||||
urgent := false
|
||||
lines := []string{}
|
||||
for i, e := range list {
|
||||
if i == 40 {
|
||||
lines = append(lines, fmt.Sprintf("and %d more", len(list)-40))
|
||||
break
|
||||
}
|
||||
key := e.Key
|
||||
if _, ok := Check(key); !ok {
|
||||
key = "(a key withheld)"
|
||||
}
|
||||
lines = append(lines, e.Severity+" "+e.What+": "+key)
|
||||
urgent = urgent || e.Severity == Urgent
|
||||
}
|
||||
m := Message{
|
||||
Title: fmt.Sprintf("HELD BACK: %d message(s) over the cap of %d an hour", len(list), CapPerHour),
|
||||
Body: strings.Join(lines, "\n") + "\nmore: conditions (through the mesh)",
|
||||
Urgent: urgent,
|
||||
}
|
||||
if ch.Ready() != nil {
|
||||
continue
|
||||
}
|
||||
_, err := ch.Send(m)
|
||||
h.noteChannel(ch.Name(), err)
|
||||
if err != nil {
|
||||
h.record(Sent{At: now, Channel: ch.Name(), Key: "(folded)", What: "fold", Outcome: "failed: " + err.Error()})
|
||||
continue
|
||||
}
|
||||
h.mu.Lock()
|
||||
h.folds[ch.Name()] = h.folds[ch.Name()][len(list):]
|
||||
h.lastFold[ch.Name()] = now
|
||||
h.mu.Unlock()
|
||||
h.record(Sent{At: now, Channel: ch.Name(), Key: "(folded)", What: fmt.Sprintf("fold of %d", len(list)), Outcome: "sent"})
|
||||
}
|
||||
}
|
||||
|
||||
// allow says whether a channel is under its cap: sends in the last hour, as recorded.
|
||||
func (h *Holder) allow(channel string) bool {
|
||||
return h.sentLastHour(channel) < CapPerHour
|
||||
}
|
||||
|
||||
func (h *Holder) sentLastHour(channel string) int {
|
||||
h.mu.Lock()
|
||||
defer h.mu.Unlock()
|
||||
since := h.Now().Add(-time.Hour)
|
||||
n := 0
|
||||
for _, s := range h.recent {
|
||||
if s.Channel == channel && s.Outcome == "sent" && s.At.After(since) && s.Key != "(folded)" {
|
||||
n++
|
||||
}
|
||||
}
|
||||
return n
|
||||
}
|
||||
|
||||
// sayOnce sends a message of the holder's own at most once per interval, on every channel ready.
|
||||
func (h *Holder) sayOnce(key string, every time.Duration, m Message) {
|
||||
now := h.Now()
|
||||
h.mu.Lock()
|
||||
if last, said := h.saidOnce[key]; said && now.Sub(last) < every {
|
||||
h.mu.Unlock()
|
||||
return
|
||||
}
|
||||
h.saidOnce[key] = now
|
||||
h.mu.Unlock()
|
||||
r := &Record{Key: key, Severity: Warning, Raised: now, Sent: map[string]string{}}
|
||||
for _, ch := range h.channels() {
|
||||
if ch.Ready() != nil {
|
||||
continue
|
||||
}
|
||||
if h.sendOn(ch, r, "notice", m) {
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func (h *Holder) record(s Sent) {
|
||||
h.mu.Lock()
|
||||
h.recent = append(h.recent, s)
|
||||
|
||||
@@ -15,8 +15,11 @@ type fakeChannel struct {
|
||||
sends []Message
|
||||
edits map[string]Message
|
||||
n int
|
||||
silent bool // edits notify nobody, as Telegram's
|
||||
}
|
||||
|
||||
func (f *fakeChannel) SilentEdit() bool { return f.silent }
|
||||
|
||||
func (f *fakeChannel) Name() string { return f.name }
|
||||
func (f *fakeChannel) Ready() error { return f.notReady }
|
||||
func (f *fakeChannel) CanEdit() bool { return true }
|
||||
@@ -72,10 +75,16 @@ func cond(key, sev, summary string) Condition {
|
||||
return Condition{Key: key, Scope: parts[0], ID: parts[1], Kind: parts[len(parts)-1], Severity: sev, Summary: summary}
|
||||
}
|
||||
|
||||
// settle lets a burst window pass and sends what is due.
|
||||
func settle(h *Holder, c *clock) {
|
||||
c.pass(BurstWindow)
|
||||
h.Flush()
|
||||
}
|
||||
|
||||
func newHolder(t *testing.T) (*Holder, *fakeChannel, *fakeChannel, *clock, *memStore) {
|
||||
t.Helper()
|
||||
c := start()
|
||||
tg, dt := &fakeChannel{name: "telegram"}, &fakeChannel{name: "desktop"}
|
||||
tg, dt := &fakeChannel{name: "telegram", silent: true}, &fakeChannel{name: "desktop"}
|
||||
st := &memStore{}
|
||||
var emitted []string
|
||||
h := &Holder{Telegram: tg, Desktop: dt, Store: st, Now: c.now, Logf: t.Logf,
|
||||
@@ -104,67 +113,73 @@ func TestARaisedConditionIsSentOnceByItsKey(t *testing.T) {
|
||||
}
|
||||
|
||||
func TestAWarningGoesToTheDesktopWhenASessionAnswersElseTelegram(t *testing.T) {
|
||||
h, tg, dt, _, _ := newHolder(t)
|
||||
h, tg, dt, c, _ := newHolder(t)
|
||||
h.Condition(EventRaised, cond("plan.41.stalled", Warning, "plan 41 waits on a build"))
|
||||
if len(dt.sends)+len(tg.sends) != 0 {
|
||||
t.Fatalf("a warning went out before its burst window")
|
||||
}
|
||||
settle(h, c)
|
||||
if len(dt.sends) != 1 || len(tg.sends) != 0 {
|
||||
t.Fatalf("desktop answered: desktop %d telegram %d", len(dt.sends), len(tg.sends))
|
||||
}
|
||||
dt.fail = errors.New("the account is not logged in")
|
||||
c.pass(DesktopWarningEvery)
|
||||
h.Condition(EventRaised, cond("plan.42.stalled", Warning, "plan 42 waits on a build"))
|
||||
settle(h, c) // the desktop fails: handed to Telegram
|
||||
h.Flush()
|
||||
if len(tg.sends) != 1 {
|
||||
t.Fatalf("no session: telegram %d", len(tg.sends))
|
||||
}
|
||||
dt.fail, dt.notReady = nil, errors.New("not configured")
|
||||
h.Condition(EventRaised, cond("plan.43.stalled", Warning, "plan 43 waits on a build"))
|
||||
settle(h, c)
|
||||
if len(tg.sends) != 2 {
|
||||
t.Fatalf("no desktop configured: telegram %d", len(tg.sends))
|
||||
}
|
||||
}
|
||||
|
||||
func TestTheCapHoldsTheRestAndFoldsThemIntoOneMessage(t *testing.T) {
|
||||
func TestTheCapIsALastLineAndIsSaidOnce(t *testing.T) {
|
||||
h, tg, _, c, _ := newHolder(t)
|
||||
h.Desktop = nil
|
||||
// Urgent conditions, each alone and after a quiet half minute: each is its own message — until
|
||||
// the cap, which is said once, in a message of its own, and holds the rest.
|
||||
for i := 0; i < 25; i++ {
|
||||
h.Condition(EventRaised, cond(fmt.Sprintf("machine.m%d.silent", i), Urgent, "a machine is silent"))
|
||||
c.pass(BurstWindow + time.Second)
|
||||
h.Flush()
|
||||
}
|
||||
if len(tg.sends) != CapPerHour {
|
||||
t.Fatalf("sent %d, cap %d", len(tg.sends), CapPerHour)
|
||||
if len(tg.sends) != CapPerHour+1 {
|
||||
t.Fatalf("sent %d, want the cap of %d and one saying so", len(tg.sends), CapPerHour)
|
||||
}
|
||||
if last := tg.sends[CapPerHour]; !strings.HasPrefix(last.Title, "HELD BACK: telegram is at its cap") {
|
||||
t.Fatalf("the cap was not said: %q", last.Title)
|
||||
}
|
||||
st := h.Status("listening")
|
||||
if st.Channels[0].Held != 5 || st.Channels[0].SentLastHour != CapPerHour {
|
||||
t.Fatalf("status: %+v", st.Channels[0])
|
||||
}
|
||||
// Said, not silent: the fold goes out at the first tick, and again only after ten minutes.
|
||||
h.Tick()
|
||||
if len(tg.sends) != CapPerHour+1 {
|
||||
t.Fatalf("no fold: %d sends", len(tg.sends))
|
||||
for i := 0; i < 10; i++ {
|
||||
c.pass(time.Minute)
|
||||
h.Tick()
|
||||
}
|
||||
fold := tg.sends[len(tg.sends)-1]
|
||||
if !strings.Contains(fold.Title, "HELD BACK: 5") {
|
||||
t.Fatalf("fold: %q", fold.Text())
|
||||
if len(tg.sends) != CapPerHour+1 {
|
||||
t.Fatalf("said the cap again or sent past it: %d", len(tg.sends))
|
||||
}
|
||||
// The hour frees the channel: what was held goes out as one digest.
|
||||
c.pass(time.Hour)
|
||||
h.Tick()
|
||||
if len(tg.sends) != CapPerHour+2 {
|
||||
t.Fatalf("sends %d", len(tg.sends))
|
||||
}
|
||||
d := tg.sends[len(tg.sends)-1]
|
||||
if !strings.HasPrefix(d.Title, "URGENT: 5 new urgent") {
|
||||
t.Fatalf("digest: %q", d.Title)
|
||||
}
|
||||
for i := 20; i < 25; i++ {
|
||||
if !strings.Contains(fold.Body, fmt.Sprintf("machine.m%d.silent", i)) {
|
||||
t.Fatalf("fold does not name m%d: %q", i, fold.Body)
|
||||
if !strings.Contains(d.Body, fmt.Sprintf("machine.m%d.silent", i)) {
|
||||
t.Fatalf("the digest does not name m%d: %q", i, d.Body)
|
||||
}
|
||||
}
|
||||
h.Condition(EventRaised, cond("machine.late.silent", Urgent, "a machine is silent"))
|
||||
c.pass(time.Minute)
|
||||
h.Tick()
|
||||
if len(tg.sends) != CapPerHour+1 {
|
||||
t.Fatalf("a second fold inside ten minutes")
|
||||
}
|
||||
c.pass(FoldEvery)
|
||||
h.Tick()
|
||||
if last := tg.sends[len(tg.sends)-1]; !strings.Contains(last.Body, "machine.late.silent") {
|
||||
t.Fatalf("the second fold: %q", last.Text())
|
||||
}
|
||||
// An hour on, the window is free again.
|
||||
c.pass(time.Hour)
|
||||
h.Condition(EventRaised, cond("machine.next.silent", Urgent, "a machine is silent"))
|
||||
if last := tg.sends[len(tg.sends)-1]; !strings.Contains(last.Text(), "machine.next.silent") {
|
||||
t.Fatalf("not sent after the window: %q", last.Text())
|
||||
}
|
||||
}
|
||||
|
||||
func TestAMessageCarryingAnAddressIsRefusedAndSaidWithItsWordsWithheld(t *testing.T) {
|
||||
@@ -219,6 +234,7 @@ func TestStillOpenPastItsBoundItIsSaidOnceMore(t *testing.T) {
|
||||
w.Raised = c.now()
|
||||
h.Condition(EventRaised, u)
|
||||
h.Condition(EventRaised, w)
|
||||
settle(h, c)
|
||||
c.pass(59 * time.Minute)
|
||||
h.Tick()
|
||||
if len(tg.sends) != 2 {
|
||||
@@ -236,6 +252,7 @@ func TestStillOpenPastItsBoundItIsSaidOnceMore(t *testing.T) {
|
||||
}
|
||||
c.pass(9 * time.Hour) // the warning is now 13 h old
|
||||
h.Tick()
|
||||
settle(h, c)
|
||||
if len(tg.sends) != 4 || !strings.Contains(tg.sends[3].Text(), "plan.41.stalled") {
|
||||
t.Fatalf("warning reminder: %d", len(tg.sends))
|
||||
}
|
||||
@@ -258,6 +275,7 @@ func TestClearedEditsTheFirstMessageAndReopenedWithinTenMinutesIsNotNew(t *testi
|
||||
t.Fatalf("still open")
|
||||
}
|
||||
c.pass(5 * time.Minute)
|
||||
k.At = c.now() // a reopening keeps its first raising time; the event says when it happened
|
||||
h.Condition(EventRaised, k)
|
||||
if len(tg.sends) != 1 || !strings.Contains(tg.edits["1"].Title, "open again, 2 times") {
|
||||
t.Fatalf("reopened as new: sends %d, edit %q", len(tg.sends), tg.edits["1"].Title)
|
||||
@@ -268,6 +286,7 @@ func TestClearedEditsTheFirstMessageAndReopenedWithinTenMinutesIsNotNew(t *testi
|
||||
if _, kept := st.recs[k.Key]; kept {
|
||||
t.Fatalf("a cleared message kept past the reopen window")
|
||||
}
|
||||
k.At = c.now()
|
||||
h.Condition(EventRaised, k)
|
||||
if len(tg.sends) != 2 {
|
||||
t.Fatalf("a raise after the window is a new message: %d", len(tg.sends))
|
||||
@@ -299,19 +318,30 @@ func TestASilencedConditionSendsNothing(t *testing.T) {
|
||||
}
|
||||
|
||||
func TestEscalationIsSaidOnce(t *testing.T) {
|
||||
h, tg, dt, _, _ := newHolder(t)
|
||||
h, tg, dt, c, _ := newHolder(t)
|
||||
k := cond("machine.novox.silent", Warning, "the anchor is silent")
|
||||
h.Condition(EventRaised, k)
|
||||
settle(h, c)
|
||||
k.Severity = Urgent
|
||||
h.Condition(EventChanged, k)
|
||||
h.Condition(EventChanged, k)
|
||||
settle(h, c)
|
||||
if len(dt.sends) != 2 || len(tg.sends) != 1 || !strings.HasPrefix(tg.sends[0].Title, "NOW URGENT") {
|
||||
t.Fatalf("desktop %d telegram %d", len(dt.sends), len(tg.sends))
|
||||
}
|
||||
// Escalated while the warning still waited its window: one message, the urgent one.
|
||||
k2 := cond("machine.ace.silent", Warning, "the home server is silent")
|
||||
h.Condition(EventRaised, k2)
|
||||
k2.Severity = Urgent
|
||||
h.Condition(EventChanged, k2)
|
||||
settle(h, c)
|
||||
if len(dt.sends) != 3 || !strings.HasPrefix(dt.sends[2].Title, "NOW URGENT") {
|
||||
t.Fatalf("desktop %d: %q", len(dt.sends), dt.sends[len(dt.sends)-1].Title)
|
||||
}
|
||||
}
|
||||
|
||||
func TestAChannelThatCannotSendSaysSoAndIsTriedAgain(t *testing.T) {
|
||||
h, tg, _, _, _ := newHolder(t)
|
||||
h, tg, _, c, _ := newHolder(t)
|
||||
h.Desktop = nil
|
||||
tg.fail = errors.New("telegram sendMessage: cannot connect (dial)")
|
||||
h.Condition(EventRaised, cond("machine.ace.silent", Urgent, "silent"))
|
||||
@@ -320,6 +350,11 @@ func TestAChannelThatCannotSendSaysSoAndIsTriedAgain(t *testing.T) {
|
||||
t.Fatalf("status: %+v", st)
|
||||
}
|
||||
tg.fail = nil
|
||||
h.Flush()
|
||||
if len(tg.sends) != 0 {
|
||||
t.Fatalf("tried again before a minute")
|
||||
}
|
||||
c.pass(RetryEvery)
|
||||
h.Tick()
|
||||
if len(tg.sends) != 1 || h.Status("listening").Verdict != "ok" {
|
||||
t.Fatalf("not retried: %d, %s", len(tg.sends), h.Status("listening").Verdict)
|
||||
@@ -361,20 +396,26 @@ func TestAnUnreadableEventIsToldOnceAnHour(t *testing.T) {
|
||||
h.Desktop = nil
|
||||
h.Unreadable("mesh-controller.condition-raised", errors.New("no severity"))
|
||||
h.Unreadable("mesh-controller.condition-raised", errors.New("no severity"))
|
||||
settle(h, c)
|
||||
if len(tg.sends) != 1 || h.Status("listening").Unreadable != 2 {
|
||||
t.Fatalf("sends %d", len(tg.sends))
|
||||
}
|
||||
c.pass(61 * time.Minute)
|
||||
h.Unreadable("mesh-controller.condition-raised", errors.New("no severity"))
|
||||
settle(h, c)
|
||||
if len(tg.sends) != 2 {
|
||||
t.Fatalf("not said again after an hour")
|
||||
}
|
||||
}
|
||||
|
||||
func TestNotifyIsKeptApartFromConditions(t *testing.T) {
|
||||
h, tg, _, _, _ := newHolder(t)
|
||||
h, tg, _, c, _ := newHolder(t)
|
||||
h.Desktop = nil
|
||||
out, err := h.Notify("backup.ace.failed", Warning, "last night's backup of the home server failed", "backup", false)
|
||||
if err != nil || out["held"] != true {
|
||||
t.Fatalf("%v %v", out, err)
|
||||
}
|
||||
settle(h, c)
|
||||
if err != nil || out["key"] != "notify.backup.ace.failed" || len(tg.sends) != 1 {
|
||||
t.Fatalf("%v %v %d", out, err, len(tg.sends))
|
||||
}
|
||||
|
||||
@@ -181,9 +181,25 @@ func run(h *Holder, l *listening) {
|
||||
logf("[messenger] %s cannot send: %v", ch.Name(), err)
|
||||
}
|
||||
}
|
||||
// The state first: what is open now is read from the controller before any event is taken, so
|
||||
// what the stream replays is history (novox/hq issue 271). Not reachable yet, it is asked for
|
||||
// again in the loop below, and old events meanwhile are kept as state, never said.
|
||||
h.Start()
|
||||
syncFromController(h)
|
||||
go func() {
|
||||
for range time.Tick(time.Minute) {
|
||||
h.Tick()
|
||||
tick := time.NewTicker(5 * time.Second)
|
||||
defer tick.Stop()
|
||||
lastTick, lastAsk := time.Now(), time.Now()
|
||||
for now := range tick.C {
|
||||
h.Flush()
|
||||
if now.Sub(lastTick) >= time.Minute {
|
||||
lastTick = now
|
||||
h.Tick()
|
||||
}
|
||||
if h.WantsSync() && now.Sub(lastAsk) >= time.Minute {
|
||||
lastAsk = now
|
||||
syncFromController(h)
|
||||
}
|
||||
}
|
||||
}()
|
||||
handle := func(e stdio.Envelope) error {
|
||||
@@ -212,6 +228,36 @@ func run(h *Holder, l *listening) {
|
||||
}
|
||||
}
|
||||
|
||||
// syncFromController reads the controller's open conditions and gives them to the holder.
|
||||
func syncFromController(h *Holder) {
|
||||
asked := h.Now()
|
||||
type answer struct {
|
||||
raw json.RawMessage
|
||||
err error
|
||||
}
|
||||
done := make(chan answer, 1)
|
||||
go func() {
|
||||
raw, err := stdio.Ask("seat:"+ControllerSeat+".conditions", map[string]any{})
|
||||
done <- answer{raw, err}
|
||||
}()
|
||||
var a answer
|
||||
select {
|
||||
case a = <-done:
|
||||
case <-time.After(30 * time.Second):
|
||||
a.err = errors.New("no answer in 30 s")
|
||||
}
|
||||
if a.err != nil {
|
||||
h.SyncFailed(a.err)
|
||||
return
|
||||
}
|
||||
open, err := DecodeOpen(a.raw)
|
||||
if err != nil {
|
||||
h.SyncFailed(err)
|
||||
return
|
||||
}
|
||||
h.Sync(open, asked)
|
||||
}
|
||||
|
||||
func str(description string) map[string]any {
|
||||
return map[string]any{"type": "string", "description": description}
|
||||
}
|
||||
@@ -230,10 +276,11 @@ func tools(h *Holder, l *listening) []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.",
|
||||
"reminded, held (for its burst window, the desktop's rhythm or the cap), refused, not sent yet, or " +
|
||||
"never said (known from the state or an old event: history, not news).",
|
||||
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 " +
|
||||
Description: "What was said to the operator lately, newest first — each send, digest, edit 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 }},
|
||||
@@ -274,7 +321,7 @@ func tools(h *Holder, l *listening) []stdio.Tool {
|
||||
return st, nil
|
||||
}},
|
||||
{Name: "messenger_recent",
|
||||
Description: "The recent sends — each message, edit, fold and failure on each channel, newest first — and the refusals.",
|
||||
Description: "The recent sends — each message, digest, edit 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",
|
||||
|
||||
@@ -58,7 +58,7 @@ func TestItDeclaresAndHoldsTheOperatorChannel(t *testing.T) {
|
||||
ControllerSeat + "." + EventRaised, ControllerSeat + "." + EventChanged, ControllerSeat + "." + EventCleared}) {
|
||||
t.Fatalf("consumes %v", m.Consumes)
|
||||
}
|
||||
if !reflect.DeepEqual(m.Emits, []string{"refused"}) || !reflect.DeepEqual(m.Invokes, []string{"seat:node-notifier.send"}) {
|
||||
if !reflect.DeepEqual(m.Emits, []string{"refused"}) || !reflect.DeepEqual(m.Invokes, []string{"seat:node-notifier.send", "seat:" + ControllerSeat + ".conditions"}) {
|
||||
t.Fatalf("emits %v invokes %v", m.Emits, m.Invokes)
|
||||
}
|
||||
if !reflect.DeepEqual(m.State, []string{"open", "sent"}) {
|
||||
|
||||
@@ -0,0 +1,482 @@
|
||||
package main
|
||||
|
||||
// The outbox: nothing is sent the moment it is decided. Each channel holds what is to be said for a
|
||||
// short window and sends it as one message — the message itself when it is alone, a digest when it
|
||||
// is not — so a burst of conditions is one notification, never one each (novox/hq issue 271).
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"strings"
|
||||
"time"
|
||||
)
|
||||
|
||||
// item is one thing to say on one channel: a record's message (Key, What), or a message of the
|
||||
// holder's own (Msg).
|
||||
type item struct {
|
||||
Key string
|
||||
What string
|
||||
Severity string
|
||||
At time.Time
|
||||
Msg *Message
|
||||
}
|
||||
|
||||
func (i item) urgent() bool { return i.Severity == Urgent && i.What != "cleared" }
|
||||
|
||||
// deliver decides where a record's message goes and holds it there. Urgent goes to every channel
|
||||
// that can send; a warning to the desktop when it is configured, otherwise to Telegram — and from
|
||||
// the desktop to Telegram when no session there takes it. No channel at all: it is unsent, said in
|
||||
// the status, and tried again every minute.
|
||||
func (h *Holder) deliver(r *Record, what string) {
|
||||
var to []Channel
|
||||
if r.Severity == Urgent {
|
||||
for _, ch := range h.channels() {
|
||||
if ch.Ready() == nil {
|
||||
to = append(to, ch)
|
||||
}
|
||||
}
|
||||
} else {
|
||||
switch {
|
||||
case h.Desktop != nil && h.Desktop.Ready() == nil:
|
||||
to = []Channel{h.Desktop}
|
||||
case h.Telegram != nil && h.Telegram.Ready() == nil:
|
||||
to = []Channel{h.Telegram}
|
||||
}
|
||||
}
|
||||
if len(to) == 0 {
|
||||
h.mu.Lock()
|
||||
r.Pending = what
|
||||
h.mu.Unlock()
|
||||
h.Logf("[messenger] could not send %s %s on any channel; trying again every minute: %s", what, r.Key, h.whyNot())
|
||||
return
|
||||
}
|
||||
now := h.Now()
|
||||
for _, ch := range to {
|
||||
h.enqueue(ch.Name(), item{Key: r.Key, What: what, Severity: r.Severity, At: now})
|
||||
}
|
||||
h.flush(now, false)
|
||||
}
|
||||
|
||||
// enqueue holds an item on a channel; an item for the same key replaces the one held before.
|
||||
func (h *Holder) enqueue(channel string, it item) {
|
||||
h.mu.Lock()
|
||||
defer h.mu.Unlock()
|
||||
list := h.outbox[channel]
|
||||
for i := range list {
|
||||
if list[i].Msg == nil && it.Msg == nil && list[i].Key == it.Key {
|
||||
at := list[i].At
|
||||
list[i] = it
|
||||
list[i].At = at
|
||||
return
|
||||
}
|
||||
}
|
||||
h.outbox[channel] = append(list, it)
|
||||
}
|
||||
|
||||
// unqueue drops what is held for a key, on every channel: it is over before it was said.
|
||||
func (h *Holder) unqueue(key string) int {
|
||||
h.mu.Lock()
|
||||
defer h.mu.Unlock()
|
||||
n := 0
|
||||
for ch, list := range h.outbox {
|
||||
kept := list[:0]
|
||||
for _, it := range list {
|
||||
if it.Msg == nil && it.Key == key {
|
||||
n++
|
||||
continue
|
||||
}
|
||||
kept = append(kept, it)
|
||||
}
|
||||
h.outbox[ch] = kept
|
||||
}
|
||||
return n
|
||||
}
|
||||
|
||||
func (h *Holder) queued(key string) bool {
|
||||
h.mu.Lock()
|
||||
defer h.mu.Unlock()
|
||||
for _, list := range h.outbox {
|
||||
for _, it := range list {
|
||||
if it.Msg == nil && it.Key == key {
|
||||
return true
|
||||
}
|
||||
}
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
// edit changes the first message on each channel that showed it. A channel that cannot, or whose
|
||||
// edit fails, or that said it inside a digest, is given a new message instead — held like any other.
|
||||
// silentOnly: only where an edit notifies nobody; elsewhere nothing.
|
||||
func (h *Holder) edit(r *Record, what string, silentOnly bool) {
|
||||
m := h.say(r, what)
|
||||
h.mu.Lock()
|
||||
sent := map[string]string{}
|
||||
for k, v := range r.Sent {
|
||||
sent[k] = v
|
||||
}
|
||||
h.mu.Unlock()
|
||||
now := h.Now()
|
||||
for _, ch := range h.channels() {
|
||||
id, shown := sent[ch.Name()]
|
||||
if !shown {
|
||||
continue
|
||||
}
|
||||
silent := editsSilently(ch)
|
||||
if silentOnly && (!silent || id == "") {
|
||||
continue
|
||||
}
|
||||
if id != "" && ch.CanEdit() {
|
||||
err := ch.Edit(id, m)
|
||||
h.noteChannel(ch.Name(), err)
|
||||
if err == nil {
|
||||
h.record(Sent{At: now, Channel: ch.Name(), Key: r.Key, What: what, Outcome: "edited"})
|
||||
continue
|
||||
}
|
||||
h.Logf("[messenger] could not edit the message for %s on %s (%v)", r.Key, ch.Name(), err)
|
||||
if silentOnly {
|
||||
continue
|
||||
}
|
||||
}
|
||||
h.enqueue(ch.Name(), item{Key: r.Key, What: what, Severity: r.Severity, At: now})
|
||||
}
|
||||
h.flush(now, false)
|
||||
}
|
||||
|
||||
// sayOnce holds a message of the holder's own at most once per interval, on the first channel ready.
|
||||
func (h *Holder) sayOnce(key string, every time.Duration, m Message) {
|
||||
now := h.Now()
|
||||
h.mu.Lock()
|
||||
if last, said := h.saidOnce[key]; said && now.Sub(last) < every {
|
||||
h.mu.Unlock()
|
||||
return
|
||||
}
|
||||
h.saidOnce[key] = now
|
||||
h.mu.Unlock()
|
||||
for _, ch := range h.channels() {
|
||||
if ch.Ready() != nil {
|
||||
continue
|
||||
}
|
||||
msg := m
|
||||
h.enqueue(ch.Name(), item{Key: key, What: "notice", Severity: Warning, At: now, Msg: &msg})
|
||||
h.flush(now, false)
|
||||
return
|
||||
}
|
||||
}
|
||||
|
||||
// flush sends what is due on every channel. force sends what is held whatever its window — for the
|
||||
// one message of urgent conditions found open at a reading of the state.
|
||||
func (h *Holder) flush(now time.Time, force bool) {
|
||||
for _, ch := range h.channels() {
|
||||
h.flushChannel(ch, now, force)
|
||||
}
|
||||
}
|
||||
|
||||
// due says whether a channel's held items go out now.
|
||||
func (h *Holder) due(ch Channel, list []item, now time.Time) bool {
|
||||
h.mu.Lock()
|
||||
grace := h.inGrace(now)
|
||||
lastTry, lastSend, lastWarn := h.lastTry[ch.Name()], h.lastSend[ch.Name()], h.lastWarn[ch.Name()]
|
||||
h.mu.Unlock()
|
||||
if len(list) == 0 || grace {
|
||||
return false
|
||||
}
|
||||
if !lastTry.IsZero() && now.Sub(lastTry) < RetryEvery {
|
||||
return false
|
||||
}
|
||||
oldest := list[0].At
|
||||
urgent := false
|
||||
for _, it := range list {
|
||||
if it.At.Before(oldest) {
|
||||
oldest = it.At
|
||||
}
|
||||
urgent = urgent || it.urgent()
|
||||
}
|
||||
if urgent {
|
||||
// At once when it is alone and the channel was quiet; otherwise with the burst it is in.
|
||||
alone := len(list) == 1 && (lastSend.IsZero() || now.Sub(lastSend) >= BurstWindow)
|
||||
return alone || now.Sub(oldest) >= BurstWindow
|
||||
}
|
||||
if now.Sub(oldest) < BurstWindow {
|
||||
return false
|
||||
}
|
||||
if ch == h.Desktop && !lastWarn.IsZero() && now.Sub(lastWarn) < DesktopWarningEvery {
|
||||
return false
|
||||
}
|
||||
return true
|
||||
}
|
||||
|
||||
// line is what an item says, resolved against the record now: nil when it is no longer to be said.
|
||||
func (h *Holder) resolve(it item, now time.Time) (*Record, *Message) {
|
||||
if it.Msg != nil {
|
||||
return nil, it.Msg
|
||||
}
|
||||
h.mu.Lock()
|
||||
r := h.open[it.Key]
|
||||
h.mu.Unlock()
|
||||
if r == nil || r.silenced(now) {
|
||||
return nil, nil
|
||||
}
|
||||
if (it.What == "cleared") != !r.Cleared.IsZero() {
|
||||
return nil, nil
|
||||
}
|
||||
m := h.say(r, it.What)
|
||||
return r, &m
|
||||
}
|
||||
|
||||
func (h *Holder) flushChannel(ch Channel, now time.Time, force bool) {
|
||||
name := ch.Name()
|
||||
h.mu.Lock()
|
||||
list := append([]item(nil), h.outbox[name]...)
|
||||
h.mu.Unlock()
|
||||
if len(list) == 0 || (!force && !h.due(ch, list, now)) {
|
||||
return
|
||||
}
|
||||
type entry struct {
|
||||
it item
|
||||
r *Record
|
||||
msg Message
|
||||
}
|
||||
var entries []entry
|
||||
for _, it := range list {
|
||||
r, m := h.resolve(it, now)
|
||||
if m == nil {
|
||||
continue
|
||||
}
|
||||
entries = append(entries, entry{it, r, *m})
|
||||
}
|
||||
taken := len(list)
|
||||
if len(entries) == 0 {
|
||||
h.drop(name, taken)
|
||||
return
|
||||
}
|
||||
if err := ch.Ready(); err != nil {
|
||||
h.noteChannel(name, err)
|
||||
h.failed(ch, list, taken, now, err)
|
||||
return
|
||||
}
|
||||
if !h.allow(name) {
|
||||
// The last line of defence: said once, in a message of its own, and the rest held.
|
||||
h.mu.Lock()
|
||||
for _, e := range entries {
|
||||
if e.r != nil {
|
||||
e.r.Folded = true
|
||||
}
|
||||
}
|
||||
said := h.capSaid[name]
|
||||
h.mu.Unlock()
|
||||
if said.IsZero() || now.Sub(said) >= time.Hour {
|
||||
m := Message{Title: fmt.Sprintf("HELD BACK: %s is at its cap of %d messages an hour", name, CapPerHour),
|
||||
Body: fmt.Sprintf("%d message(s) are held and go out as one when the hour allows.\n"+
|
||||
"more: conditions (through the mesh)", len(entries)), Urgent: true}
|
||||
_, err := ch.Send(m)
|
||||
h.noteChannel(name, err)
|
||||
if err == nil {
|
||||
h.mu.Lock()
|
||||
h.capSaid[name] = now
|
||||
h.mu.Unlock()
|
||||
h.record(Sent{At: now, Channel: name, Key: "(cap)", What: "cap reached", Outcome: "said once"})
|
||||
}
|
||||
}
|
||||
return
|
||||
}
|
||||
var m Message
|
||||
if len(entries) == 1 {
|
||||
m = entries[0].msg
|
||||
} else {
|
||||
items := make([]item, len(entries))
|
||||
msgs := make([]Message, len(entries))
|
||||
for i, e := range entries {
|
||||
items[i], msgs[i] = e.it, e.msg
|
||||
}
|
||||
m = digest(items, msgs)
|
||||
}
|
||||
id, err := ch.Send(m)
|
||||
h.noteChannel(name, err)
|
||||
if err != nil {
|
||||
h.failed(ch, list, taken, now, err)
|
||||
return
|
||||
}
|
||||
h.drop(name, taken)
|
||||
h.mu.Lock()
|
||||
h.lastSend[name] = now
|
||||
h.lastTry[name] = time.Time{}
|
||||
for _, e := range entries {
|
||||
if ch == h.Desktop && !e.it.urgent() {
|
||||
// The desktop's warning rhythm counts from a message that carried a warning.
|
||||
h.lastWarn[name] = now
|
||||
}
|
||||
}
|
||||
for _, e := range entries {
|
||||
r := e.r
|
||||
if r == nil {
|
||||
continue
|
||||
}
|
||||
r.Pending, r.Folded, r.History = "", false, false
|
||||
if e.it.What != "cleared" {
|
||||
if r.Sent == nil {
|
||||
r.Sent = map[string]string{}
|
||||
}
|
||||
if _, has := r.Sent[name]; !has {
|
||||
own := ""
|
||||
if len(entries) == 1 {
|
||||
own = id
|
||||
}
|
||||
r.Sent[name] = own
|
||||
}
|
||||
if r.FirstSent.IsZero() {
|
||||
r.FirstSent = now
|
||||
}
|
||||
}
|
||||
}
|
||||
h.mu.Unlock()
|
||||
key, what := "(digest)", fmt.Sprintf("digest of %d", len(entries))
|
||||
if len(entries) == 1 {
|
||||
key, what = entries[0].it.Key, entries[0].it.What
|
||||
}
|
||||
h.record(Sent{At: now, Channel: name, Key: key, What: what, Outcome: "sent"})
|
||||
for _, e := range entries {
|
||||
if len(entries) > 1 {
|
||||
h.record(Sent{At: now, Channel: name, Key: e.it.Key, What: e.it.What, Outcome: "in digest"})
|
||||
}
|
||||
if e.r != nil {
|
||||
h.persist(e.r)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// drop removes the first n held items of a channel: those this flush took.
|
||||
func (h *Holder) drop(channel string, n int) {
|
||||
h.mu.Lock()
|
||||
defer h.mu.Unlock()
|
||||
list := h.outbox[channel]
|
||||
if n > len(list) {
|
||||
n = len(list)
|
||||
}
|
||||
h.outbox[channel] = append([]item(nil), list[n:]...)
|
||||
}
|
||||
|
||||
// failed keeps what a channel could not send, to be tried again in a minute. A warning the desktop
|
||||
// could not show — no session there — goes to Telegram instead.
|
||||
func (h *Holder) failed(ch Channel, list []item, taken int, now time.Time, err error) {
|
||||
name := ch.Name()
|
||||
h.mu.Lock()
|
||||
h.lastTry[name] = now
|
||||
h.mu.Unlock()
|
||||
h.record(Sent{At: now, Channel: name, Key: "(held)", What: fmt.Sprintf("%d held", taken), Outcome: "failed: " + err.Error()})
|
||||
if ch == h.Desktop && h.Telegram != nil && h.Telegram.Ready() == nil {
|
||||
var keep []item
|
||||
for _, it := range list[:taken] {
|
||||
if it.Severity == Urgent {
|
||||
keep = append(keep, it)
|
||||
continue
|
||||
}
|
||||
h.enqueue(h.Telegram.Name(), it)
|
||||
}
|
||||
h.mu.Lock()
|
||||
h.outbox[name] = append(keep, h.outbox[name][taken:]...)
|
||||
h.mu.Unlock()
|
||||
}
|
||||
h.mu.Lock()
|
||||
for _, it := range list[:taken] {
|
||||
if r := h.open[it.Key]; r != nil && it.Msg == nil && r.Cleared.IsZero() {
|
||||
r.Pending = it.What
|
||||
}
|
||||
}
|
||||
h.mu.Unlock()
|
||||
}
|
||||
|
||||
// digest is several messages as one: what they are, counted, then each in a line.
|
||||
func digest(items []item, msgs []Message) Message {
|
||||
var newU, newW, still, cleared, notices int
|
||||
urgent, quiet := false, true
|
||||
for _, it := range items {
|
||||
switch {
|
||||
case it.Msg != nil:
|
||||
notices++
|
||||
case it.What == "cleared":
|
||||
cleared++
|
||||
case it.What == "reminder":
|
||||
still++
|
||||
case it.Severity == Urgent:
|
||||
newU++
|
||||
default:
|
||||
newW++
|
||||
}
|
||||
urgent = urgent || it.urgent()
|
||||
quiet = quiet && it.What == "cleared"
|
||||
}
|
||||
var parts []string
|
||||
add := func(n int, one, many string) {
|
||||
switch {
|
||||
case n == 1:
|
||||
parts = append(parts, "1 "+one)
|
||||
case n > 1:
|
||||
parts = append(parts, fmt.Sprintf("%d %s", n, many))
|
||||
}
|
||||
}
|
||||
add(newU, "new urgent", "new urgent")
|
||||
add(newW, "new warning", "new warnings")
|
||||
add(still, "still open", "still open")
|
||||
add(cleared, "cleared", "cleared")
|
||||
add(notices, "notice", "notices")
|
||||
prefix := "WARNING"
|
||||
if urgent {
|
||||
prefix = "URGENT"
|
||||
} else if quiet {
|
||||
prefix = "CLEARED"
|
||||
}
|
||||
var lines []string
|
||||
for i, m := range msgs {
|
||||
if i == DigestLines {
|
||||
lines = append(lines, fmt.Sprintf("and %d more", len(msgs)-DigestLines))
|
||||
break
|
||||
}
|
||||
line := m.Title
|
||||
if key := items[i].Key; items[i].Msg == nil {
|
||||
if _, ok := Check(key); ok {
|
||||
line += " [" + key + "]"
|
||||
}
|
||||
}
|
||||
lines = append(lines, "- "+line)
|
||||
}
|
||||
lines = append(lines, "more: conditions (through the mesh)")
|
||||
return Message{Title: prefix + ": " + strings.Join(parts, ", "), Body: strings.Join(lines, "\n"), Urgent: urgent, Quiet: quiet}
|
||||
}
|
||||
|
||||
// sendNow sends one message on one channel at once, under the cap: for what was asked for (a test).
|
||||
func (h *Holder) sendNow(ch Channel, key, what string, m Message) bool {
|
||||
if err := ch.Ready(); err != nil {
|
||||
h.noteChannel(ch.Name(), err)
|
||||
return false
|
||||
}
|
||||
_, err := ch.Send(m)
|
||||
h.noteChannel(ch.Name(), err)
|
||||
now := h.Now()
|
||||
if err != nil {
|
||||
h.record(Sent{At: now, Channel: ch.Name(), Key: key, What: what, Outcome: "failed: " + err.Error()})
|
||||
return false
|
||||
}
|
||||
h.mu.Lock()
|
||||
h.lastSend[ch.Name()] = now
|
||||
h.mu.Unlock()
|
||||
h.record(Sent{At: now, Channel: ch.Name(), Key: key, What: what, Outcome: "sent"})
|
||||
return true
|
||||
}
|
||||
|
||||
// allow says whether a channel is under its cap: messages sent in the last hour, as recorded.
|
||||
func (h *Holder) allow(channel string) bool {
|
||||
return h.sentLastHour(channel) < CapPerHour
|
||||
}
|
||||
|
||||
func (h *Holder) sentLastHour(channel string) int {
|
||||
h.mu.Lock()
|
||||
defer h.mu.Unlock()
|
||||
since := h.Now().Add(-time.Hour)
|
||||
n := 0
|
||||
for _, s := range h.recent {
|
||||
if s.Channel == channel && s.Outcome == "sent" && s.At.After(since) {
|
||||
n++
|
||||
}
|
||||
}
|
||||
return n
|
||||
}
|
||||
@@ -14,7 +14,7 @@ type ChannelStatus struct {
|
||||
LastSent string `json:"last_sent,omitempty"`
|
||||
SentLastHour int `json:"sent_last_hour"`
|
||||
Cap int `json:"cap_per_hour"`
|
||||
Held int `json:"held_by_cap"`
|
||||
Held int `json:"held"` // waiting for their window, the digest or the cap
|
||||
}
|
||||
|
||||
// Status is the holder's whole account of itself.
|
||||
@@ -30,6 +30,11 @@ type Status struct {
|
||||
Heard map[string]int `json:"events_heard_since_start"`
|
||||
LastHeard string `json:"last_event_heard,omitempty"`
|
||||
Listening string `json:"listening"`
|
||||
StateRead string `json:"state_read,omitempty"` // when the open conditions were last read from the controller
|
||||
StateUnread string `json:"state_unreadable,omitempty"` // why they could not be, lately
|
||||
History int `json:"events_already_in_state"` // older than the reading: not acted on
|
||||
Stale int `json:"events_taken_as_state"` // older than the freshness bound: recorded, never said
|
||||
Grace bool `json:"in_start_grace,omitempty"`
|
||||
StateProblem string `json:"state_problem,omitempty"`
|
||||
Rules []string `json:"rules"`
|
||||
}
|
||||
@@ -50,7 +55,7 @@ func (h *Holder) Status(listening string) Status {
|
||||
if t, ok := h.chanOK[ch.Name()]; ok {
|
||||
cs.LastSent = t.UTC().Format(time.RFC3339)
|
||||
}
|
||||
cs.Held = len(h.folds[ch.Name()])
|
||||
cs.Held = len(h.outbox[ch.Name()])
|
||||
h.mu.Unlock()
|
||||
if cs.LastError != "" {
|
||||
cs.CanSend = false
|
||||
@@ -81,13 +86,24 @@ func (h *Holder) Status(listening string) Status {
|
||||
st.LastHeard = h.lastHeard.UTC().Format(time.RFC3339)
|
||||
}
|
||||
st.StateProblem = h.storeErr
|
||||
if !h.syncedAt.IsZero() {
|
||||
st.StateRead = h.syncedAt.UTC().Format(time.RFC3339)
|
||||
}
|
||||
st.StateUnread = h.syncErr
|
||||
st.History, st.Stale = h.history, h.stale
|
||||
st.Grace = h.inGrace(now)
|
||||
h.mu.Unlock()
|
||||
sort.Strings(st.Unsent)
|
||||
st.Listening = listening
|
||||
st.Rules = []string{
|
||||
"history is never news: what is open is read from the controller on start, every 10 min and when an old event arrives; an event older than that reading is not acted on, one older than 10 min by its own time is state, never a message",
|
||||
"a clearing of something never said says nothing; an urgent condition open and never said is said, all of them in one message",
|
||||
"nothing is sent in the first minute after start but that one message",
|
||||
"urgent: Telegram and the desktop; warning: the desktop where a session answers, otherwise Telegram",
|
||||
"bursts are one message: what arrives within 30 s goes out together, a digest when it is several; urgent goes at once when the channel was quiet for 30 s",
|
||||
"the desktop is gentle: warnings at most one message every 15 min",
|
||||
"deduplicated by the condition's key; reminded once after 1 h (urgent) or 12 h (warning); cleared by editing the first message",
|
||||
"at most 20 messages an hour per channel; the rest folded into one message, at most every 10 min",
|
||||
"at most 20 messages an hour per channel, a last line of defence: reached, it is said once and the rest held for one digest",
|
||||
"refused: anything carrying an address, a path or a secret; channel-refused is sent with the words withheld",
|
||||
"silenced conditions send nothing; silence is `conditions silence` through the mesh",
|
||||
}
|
||||
@@ -127,7 +143,8 @@ type OpenMessage struct {
|
||||
Silenced string `json:"silenced_until,omitempty"`
|
||||
Reminded bool `json:"reminded"`
|
||||
Unsent string `json:"unsent,omitempty"`
|
||||
Held bool `json:"held_by_cap,omitempty"`
|
||||
Held bool `json:"held,omitempty"`
|
||||
History bool `json:"never_said,omitempty"`
|
||||
Refused string `json:"refused,omitempty"`
|
||||
Count int `json:"times_opened"`
|
||||
}
|
||||
@@ -137,6 +154,12 @@ func (h *Holder) Open() []OpenMessage {
|
||||
now := h.Now()
|
||||
h.mu.Lock()
|
||||
defer h.mu.Unlock()
|
||||
queued := map[string]bool{}
|
||||
for _, list := range h.outbox {
|
||||
for _, it := range list {
|
||||
queued[it.Key] = true
|
||||
}
|
||||
}
|
||||
var out []OpenMessage
|
||||
for _, r := range h.open {
|
||||
if !r.Cleared.IsZero() {
|
||||
@@ -144,7 +167,7 @@ func (h *Holder) Open() []OpenMessage {
|
||||
}
|
||||
o := OpenMessage{Key: r.Key, Severity: r.Severity, Summary: r.Summary, Subject: r.Subject,
|
||||
Since: r.Raised.UTC().Format(time.RFC3339), Origin: r.Origin, Reminded: r.Reminded,
|
||||
Unsent: r.Pending, Held: r.Folded, Refused: r.Refused, Count: r.Count}
|
||||
Unsent: r.Pending, Held: r.Folded || queued[r.Key], History: r.History, Refused: r.Refused, Count: r.Count}
|
||||
for ch := range r.Sent {
|
||||
o.SentOn = append(o.SentOn, ch)
|
||||
}
|
||||
@@ -190,7 +213,6 @@ func (h *Holder) Test(channel, text string) map[string]string {
|
||||
out["refused"] = "it carried " + refusal.String() + "; nothing was sent"
|
||||
return out
|
||||
}
|
||||
r := &Record{Key: "messenger.test", Severity: Warning, Raised: h.Now(), Sent: map[string]string{}}
|
||||
m := Message{Title: "TEST: " + text, Body: "nothing is wrong; this was asked for"}
|
||||
for _, ch := range h.channels() {
|
||||
if channel != "" && channel != "both" && channel != ch.Name() {
|
||||
@@ -204,7 +226,7 @@ func (h *Holder) Test(channel, text string) map[string]string {
|
||||
out[ch.Name()] = "not sent: the channel is at its cap of 20 an hour"
|
||||
continue
|
||||
}
|
||||
if h.sendOn(ch, r, "test", m) {
|
||||
if h.sendNow(ch, "messenger.test", "test", m) {
|
||||
out[ch.Name()] = "sent"
|
||||
} else {
|
||||
h.mu.Lock()
|
||||
@@ -231,7 +253,7 @@ func (h *Holder) Notify(key, severity, summary, subject string, clear bool) (map
|
||||
h.init()
|
||||
h.mu.Unlock()
|
||||
if clear {
|
||||
h.cleared(key)
|
||||
h.cleared(key, false)
|
||||
return map[string]any{"key": key, "cleared": true}, nil
|
||||
}
|
||||
switch severity {
|
||||
@@ -245,7 +267,7 @@ func (h *Holder) Notify(key, severity, summary, subject string, clear bool) (map
|
||||
return nil, errorf("summary is required: one line in the mesh's words")
|
||||
}
|
||||
h.raised(Record{Key: key, Kind: "notice", Subject: subject, Severity: severity, Summary: summary,
|
||||
Origin: "notify", More: "operator-channel.open"})
|
||||
Origin: "notify", More: "operator-channel.open"}, false)
|
||||
h.mu.Lock()
|
||||
r := h.open[key]
|
||||
var sent []string
|
||||
@@ -258,5 +280,6 @@ func (h *Holder) Notify(key, severity, summary, subject string, clear bool) (map
|
||||
}
|
||||
h.mu.Unlock()
|
||||
sort.Strings(sent)
|
||||
return map[string]any{"key": key, "sent_on": sent, "refused": refused, "unsent": pending}, nil
|
||||
held := h.queued(key)
|
||||
return map[string]any{"key": key, "sent_on": sent, "held": held, "refused": refused, "unsent": pending}, nil
|
||||
}
|
||||
|
||||
@@ -0,0 +1,170 @@
|
||||
package main
|
||||
|
||||
// Reading the state (novox/hq issue 271): what is open now, from the controller's `conditions`
|
||||
// verb, is what this holder believes — not the events it happened to be handed. Read on start, every
|
||||
// ResyncEvery, and at once when an old event arrives (a replay, a recreated consumer, the holder back
|
||||
// from away). Everything before the reading is in it, so no event older than it is acted on.
|
||||
|
||||
import (
|
||||
"sort"
|
||||
"time"
|
||||
)
|
||||
|
||||
// SyncResult is what a reading of the state changed, for the log and the status.
|
||||
type SyncResult struct {
|
||||
Open int `json:"open"`
|
||||
Said []string `json:"said,omitempty"` // urgent and never said: said now, in one message
|
||||
News []string `json:"news,omitempty"` // a warning raised within FreshFor and not heard: said like any other
|
||||
State []string `json:"as_state,omitempty"` // open, older, never said: kept, not said
|
||||
Gone []string `json:"gone,omitempty"` // open here, not in the controller's state: ended, quietly
|
||||
Skipped []string `json:"newer_here,omitempty"`
|
||||
}
|
||||
|
||||
// Sync takes the controller's open conditions, read by a request made at `asked`.
|
||||
func (h *Holder) Sync(open []Condition, asked time.Time) SyncResult {
|
||||
h.work.Lock()
|
||||
defer h.work.Unlock()
|
||||
h.mu.Lock()
|
||||
h.init()
|
||||
now := h.Now()
|
||||
h.mu.Unlock()
|
||||
res := SyncResult{Open: len(open)}
|
||||
in := map[string]bool{}
|
||||
for _, c := range open {
|
||||
in[c.Key] = true
|
||||
}
|
||||
// What is open here and no longer there ended while this holder was not told: forgotten without
|
||||
// a message — said only where an edit notifies nobody. What this holder heard of after the
|
||||
// reading was asked for is newer than it, and stands.
|
||||
h.mu.Lock()
|
||||
var gone []*Record
|
||||
for _, r := range h.open {
|
||||
if r.Origin == "condition" && r.Cleared.IsZero() && !in[r.Key] && !r.Seen.After(asked) {
|
||||
gone = append(gone, r)
|
||||
}
|
||||
}
|
||||
h.mu.Unlock()
|
||||
sort.Slice(gone, func(i, j int) bool { return gone[i].Key < gone[j].Key })
|
||||
for _, r := range gone {
|
||||
res.Gone = append(res.Gone, r.Key)
|
||||
h.cleared(r.Key, true)
|
||||
}
|
||||
var urgent []*Record
|
||||
var fresh []*Record
|
||||
for _, c := range open {
|
||||
rec := Record{
|
||||
Key: c.Key, Kind: c.Kind, Subject: c.SubjectWords(), Severity: c.Severity, Summary: c.Summary,
|
||||
Origin: "condition", More: "conditions show " + c.Key, Raised: c.Raised, SilencedTill: c.SilencedTill,
|
||||
}
|
||||
h.mu.Lock()
|
||||
old := h.open[c.Key]
|
||||
var r *Record
|
||||
known := old != nil
|
||||
switch {
|
||||
case old != nil && old.Seen.After(asked):
|
||||
h.mu.Unlock()
|
||||
res.Skipped = append(res.Skipped, c.Key)
|
||||
continue
|
||||
case old != nil:
|
||||
if !old.Cleared.IsZero() {
|
||||
// Ended here, open there: back to open, quietly — the reading is not an event.
|
||||
old.Cleared = time.Time{}
|
||||
old.Count++
|
||||
}
|
||||
old.Summary, old.Subject, old.Kind, old.SilencedTill = rec.Summary, rec.Subject, rec.Kind, rec.SilencedTill
|
||||
if rec.Severity != "" {
|
||||
old.Severity = rec.Severity
|
||||
}
|
||||
if !rec.Raised.IsZero() {
|
||||
old.Raised = rec.Raised
|
||||
}
|
||||
r = old
|
||||
default:
|
||||
n := rec
|
||||
n.Count = 1
|
||||
if n.Raised.IsZero() {
|
||||
n.Raised = now
|
||||
}
|
||||
n.Sent = map[string]string{}
|
||||
n.History = true
|
||||
h.open[n.Key] = &n
|
||||
r = &n
|
||||
}
|
||||
r.Seen = now
|
||||
told := r.told()
|
||||
silenced := r.silenced(now)
|
||||
h.mu.Unlock()
|
||||
switch {
|
||||
case told || silenced || h.queued(r.Key):
|
||||
case r.Severity == Urgent:
|
||||
urgent = append(urgent, r)
|
||||
case !known && now.Sub(r.Raised) <= FreshFor:
|
||||
fresh = append(fresh, r)
|
||||
default:
|
||||
res.State = append(res.State, r.Key)
|
||||
}
|
||||
h.persist(r)
|
||||
}
|
||||
// Every urgent condition open and never said: one message, now, whatever the grace.
|
||||
if len(urgent) > 0 {
|
||||
var ready []Channel
|
||||
for _, ch := range h.channels() {
|
||||
if ch.Ready() == nil {
|
||||
ready = append(ready, ch)
|
||||
}
|
||||
}
|
||||
for _, r := range urgent {
|
||||
res.Said = append(res.Said, r.Key)
|
||||
if len(ready) == 0 {
|
||||
h.mu.Lock()
|
||||
r.Pending = "raised"
|
||||
h.mu.Unlock()
|
||||
continue
|
||||
}
|
||||
for _, ch := range ready {
|
||||
h.enqueue(ch.Name(), item{Key: r.Key, What: "raised", Severity: Urgent, At: now})
|
||||
}
|
||||
}
|
||||
for _, ch := range ready {
|
||||
h.flushChannel(ch, now, true)
|
||||
}
|
||||
}
|
||||
for _, r := range fresh {
|
||||
res.News = append(res.News, r.Key)
|
||||
h.mu.Lock()
|
||||
r.History = false
|
||||
h.mu.Unlock()
|
||||
h.deliver(r, "raised")
|
||||
}
|
||||
h.mu.Lock()
|
||||
h.syncedAt = asked
|
||||
h.syncErr = ""
|
||||
h.wantSync = false
|
||||
h.mu.Unlock()
|
||||
if len(res.Said)+len(res.News)+len(res.Gone) > 0 || len(res.State) > 0 {
|
||||
h.Logf("[messenger] read the open conditions: %d open; said %d urgent never said, %d new, kept %d as state, %d ended meanwhile",
|
||||
res.Open, len(res.Said), len(res.News), len(res.State), len(res.Gone))
|
||||
}
|
||||
return res
|
||||
}
|
||||
|
||||
// SyncFailed notes that the state could not be read; it is asked for again.
|
||||
func (h *Holder) SyncFailed(err error) {
|
||||
h.mu.Lock()
|
||||
defer h.mu.Unlock()
|
||||
h.init()
|
||||
first := h.syncErr == ""
|
||||
h.syncErr = err.Error()
|
||||
if first {
|
||||
h.Logf("[messenger] cannot read the open conditions from the controller (%v): old events are kept as state, never said", err)
|
||||
}
|
||||
}
|
||||
|
||||
// WantsSync says whether the state should be read again: never read, read too long ago, or an old
|
||||
// event arrived since.
|
||||
func (h *Holder) WantsSync() bool {
|
||||
h.mu.Lock()
|
||||
defer h.mu.Unlock()
|
||||
h.init()
|
||||
return h.wantSync || h.syncedAt.IsZero() || h.Now().Sub(h.syncedAt) >= ResyncEvery
|
||||
}
|
||||
@@ -128,6 +128,9 @@ func (t *Telegram) Edit(id string, m Message) error {
|
||||
// CanEdit: Telegram edits a message in place.
|
||||
func (t *Telegram) CanEdit() bool { return true }
|
||||
|
||||
// SilentEdit: an edited Telegram message notifies nobody.
|
||||
func (t *Telegram) SilentEdit() bool { return true }
|
||||
|
||||
// Who asks the bot API who it is: the check that the token works, with no message sent.
|
||||
func (t *Telegram) Who() (string, error) {
|
||||
token, _, err := t.ready()
|
||||
|
||||
File diff suppressed because it is too large
Load Diff
Binary file not shown.
@@ -33,7 +33,8 @@
|
||||
"refused"
|
||||
],
|
||||
"invokes": [
|
||||
"seat:node-notifier.send"
|
||||
"seat:node-notifier.send",
|
||||
"seat:mesh-controller.conditions"
|
||||
],
|
||||
"state": [
|
||||
"open",
|
||||
|
||||
Reference in New Issue
Block a user