mesh/merge-gate pass: builds build-agent, mesh-controller, route-proxy → ace, g14, novox, shanks; no bus step; every machine composes with the change as it…
mesh/repo-check pass: its merge-check.sh passed
mesh/delivery delivered
mesh/delivery-group group feat/plain-notifications delivered: every member is delivered
The operator could not tell from a notification whether to act, and was told to have an agent do it. Each condition now says "Nothing for you to do." or "Needs you:" with one thing they can do themselves, and carries the actions the operator channel performs when chosen (hq ADR 0253).
184 lines
6.6 KiB
Go
184 lines
6.6 KiB
Go
package main
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"strings"
|
|
"time"
|
|
|
|
"github.com/nats-io/nats.go"
|
|
|
|
"github.com/novox/mesh-controller/internal/broker"
|
|
"github.com/novox/mesh-controller/internal/catalogue"
|
|
"github.com/novox/mesh-controller/internal/conditions"
|
|
"github.com/novox/mesh-controller/internal/link"
|
|
)
|
|
|
|
// A delivery held past its bound, said by the controller and healed by H2 from mesh-delivery's own table
|
|
// (novox/hq ADR 0239 decision 9, to-be 47 Phase B). The conditions are the controller's, one owner: the
|
|
// self-check reads the delivery owner's `stalled` (probe D14) and raises `delivery.<id>.stalled`; healer
|
|
// H2 calls the owner's `close`, which takes only the transition the table names for that state, and only
|
|
// the next observation says whether it worked (ADR 0231).
|
|
|
|
// probeDeliveriesID is the probe that reads the delivery's owner.
|
|
const probeDeliveriesID = "D14"
|
|
|
|
// kindDeliveryStalled is what a delivery held past its bound raises: H2's kind, in the delivery scope.
|
|
const kindDeliveryStalled = "stalled"
|
|
|
|
// deliveryOwnerWithin is how long the owner is given to answer.
|
|
var deliveryOwnerWithin = 10 * time.Second
|
|
|
|
// stalledLine is one delivery past its bound, as mesh-delivery's `stalled` says it.
|
|
type stalledLine struct {
|
|
ID string `json:"id"`
|
|
State string `json:"state"`
|
|
For string `json:"for"`
|
|
Bound string `json:"bound"`
|
|
H2 string `json:"h2"`
|
|
Says string `json:"says"`
|
|
}
|
|
|
|
// operatorsOnly is whether the table leaves H2 nothing to do for the line: the state is the operator's.
|
|
func (l stalledLine) operatorsOnly() bool { return l.H2 == "" || strings.HasPrefix(l.H2, "none") }
|
|
|
|
// askDeliveryOwner asks the holder of the mesh-delivery seat one of the verbs the controller is granted;
|
|
// a seam a test replaces. The holder's own refusal is an error naming it.
|
|
var askDeliveryOwner = func(ctx context.Context, conn *nats.Conn, verb string, args map[string]any) (json.RawMessage, error) {
|
|
granted := false
|
|
for _, v := range broker.VerbsTheControllerAsksTheDeliveryOwner {
|
|
granted = granted || v.Verb == verb
|
|
}
|
|
if !granted {
|
|
return nil, fmt.Errorf("the controller asks %s.%s, which its grant does not name", catalogue.DeliverySeat, verb)
|
|
}
|
|
if conn == nil {
|
|
return nil, errors.New("this controller is not on the bus")
|
|
}
|
|
answer, err := link.AskMeshSeatTool(ctx, conn, catalogue.DeliverySeat, verb, args, deliveryOwnerWithin)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
if answer.Error != "" {
|
|
return nil, fmt.Errorf("%s.%s refused: %s", catalogue.DeliverySeat, verb, answer.Error)
|
|
}
|
|
return unwrapToolResult(answer.Result), nil
|
|
}
|
|
|
|
// unwrapToolResult is a tool's answer whatever the runtime wrapped it in: the JSON itself, or the
|
|
// protocol's content list holding it as text.
|
|
func unwrapToolResult(raw json.RawMessage) json.RawMessage {
|
|
var wrapped struct {
|
|
Content []struct {
|
|
Text string `json:"text"`
|
|
} `json:"content"`
|
|
}
|
|
if json.Unmarshal(raw, &wrapped) == nil && len(wrapped.Content) > 0 && json.Valid([]byte(wrapped.Content[0].Text)) {
|
|
return json.RawMessage(wrapped.Content[0].Text)
|
|
}
|
|
return raw
|
|
}
|
|
|
|
// deliveriesStalled is what the owner says is held past its bound; nothing when no holder is on record,
|
|
// or none answers (D3 says that one).
|
|
func deliveriesStalled(ctx context.Context, conn *nats.Conn, held bool) ([]stalledLine, error) {
|
|
if !held {
|
|
return nil, nil
|
|
}
|
|
raw, err := askDeliveryOwner(ctx, conn, "stalled", map[string]any{})
|
|
if errors.Is(err, link.ErrNothingServes) {
|
|
return nil, nil // the holder not running is D3's holder-silent, said once there
|
|
}
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
var lines []stalledLine
|
|
if err := json.Unmarshal(raw, &lines); err != nil {
|
|
return nil, fmt.Errorf("%s.stalled answered something unreadable: %w", catalogue.DeliverySeat, err)
|
|
}
|
|
return lines, nil
|
|
}
|
|
|
|
// stalledObservations are the conditions of what is stalled.
|
|
func stalledObservations(lines []stalledLine) []conditions.Observation {
|
|
out := make([]conditions.Observation, 0, len(lines))
|
|
for _, l := range lines {
|
|
o := conditions.Observation{Scope: conditions.ScopeDelivery, ID: l.ID, Kind: kindDeliveryStalled,
|
|
Severity: conditions.Warning,
|
|
Summary: fmt.Sprintf("the delivery %s has been %s for %s, past its bound of %s (%s): healer H2 may %s — "+
|
|
"`mesh-delivery.show %s`", l.ID, l.State, l.For, l.Bound, l.Says, l.H2, l.ID),
|
|
Said: fmt.Sprintf("%s for %s", l.State, l.For)}
|
|
if l.operatorsOnly() {
|
|
o.Resolver = conditions.ResolverOperator
|
|
}
|
|
o.Headline, o.Explanation, o.Resolved, o.Needs, o.Actions = stalledWords(l, o)
|
|
out = append(out, o)
|
|
}
|
|
return out
|
|
}
|
|
|
|
// probeDeliveries is D14: every delivery held past its state's bound is said.
|
|
func probeDeliveries(ctx context.Context, d *doctor) ([]conditions.Observation, error) {
|
|
entries, err := d.open.inventory.Catalogued(ctx)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
var conn *nats.Conn
|
|
if d.js != nil {
|
|
conn = d.js.Conn()
|
|
}
|
|
lines, err := deliveriesStalled(ctx, conn, deliverySeatHeld(entries))
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
return stalledObservations(lines), nil
|
|
}
|
|
|
|
// --- H2, for a delivery ---------------------------------------------------------------------------------
|
|
|
|
// connOf is the bus a healer asks over.
|
|
func (h *healing) connOf() *nats.Conn {
|
|
if h.js == nil {
|
|
return nil
|
|
}
|
|
return h.js.Conn()
|
|
}
|
|
|
|
// appliesToAStalledDelivery is H2's for a delivery: the owner still lists it, with a transition the table
|
|
// lets H2 take; never one whose state is the operator's.
|
|
func appliesToAStalledDelivery(ctx context.Context, h *healing, c conditions.Condition) (string, bool, string, error) {
|
|
lines, err := deliveriesStalled(ctx, h.connOf(), true)
|
|
if err != nil {
|
|
return "", false, "", err
|
|
}
|
|
for _, l := range lines {
|
|
if l.ID != c.Subject.ID {
|
|
continue
|
|
}
|
|
if l.operatorsOnly() {
|
|
return "", false, "the delivery is " + l.State + ", a state the table leaves to the operator", nil
|
|
}
|
|
return c.Key, true, "", nil
|
|
}
|
|
return "", false, "the delivery's owner no longer lists it as stalled: its condition clears on the next look", nil
|
|
}
|
|
|
|
// repairDelivery is H2 for a delivery: the owner's `close`, which reads the walk again and takes only the
|
|
// transition its table names. A refusal is no repair, said; the next observation says whether it worked.
|
|
func repairDelivery(ctx context.Context, h *healing, c conditions.Condition) (string, string, error) {
|
|
raw, err := askDeliveryOwner(ctx, h.connOf(), "close", map[string]any{"id": c.Subject.ID, "why": c.Key})
|
|
if err != nil {
|
|
if errors.Is(err, link.ErrNothingServes) {
|
|
return "", "", err
|
|
}
|
|
return "closed nothing", oneLine(err.Error()), nil
|
|
}
|
|
var said string
|
|
if json.Unmarshal(raw, &said) != nil {
|
|
said = string(raw)
|
|
}
|
|
return "asked " + catalogue.DeliverySeat + " to close " + c.Subject.ID, said, nil
|
|
}
|