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 fail: its merge-check.sh failed: FAIL
mesh/delivery delivered
mesh/delivery-group group feat/mesh-delivery-waits-said delivered: every member is delivered
A walk mesh-delivery never lets go waited for ever with nothing open: S16 says it at 30 minutes, urgent at 4 hours, naming plans go. Phase B: probe D14 reads the delivery owner's stalled and raises delivery.<id>.stalled, and H2 takes the table's transition through its close.
183 lines
6.6 KiB
Go
183 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
|
|
}
|
|
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
|
|
}
|