Merge pull request 'messenger: history is never news — read the open conditions first, coalesce bursts (hq issue 271)' (#88) from fix/messenger-history-is-never-news into main
This commit was merged in pull request #88.
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