The gate judged a module by what the mesh saw from outside, so a container that crash-looped after it applied passed it. Each machine's node-engine now states the health of every long-running resource it runs; the controller keeps the newest statement per machine, raises module.<module>.<machine>.unhealthy on the second statement in a row, clears it on the first that does not say it, and the gate passes a module only when every long-running resource of it is stated healthy since the send. An engine that states nothing is judged as before.
224 lines
8.4 KiB
Go
224 lines
8.4 KiB
Go
package main
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"fmt"
|
|
"os"
|
|
"sync"
|
|
"time"
|
|
|
|
"github.com/novox/mesh-controller/internal/link"
|
|
)
|
|
|
|
// `status` answered from a summary the serving controller keeps current (novox/hq to-be 45 Phase 0,
|
|
// §4's D9 and §8's health of the controller).
|
|
//
|
|
// **Asked, it was composed: every machine resolved, twice, while its caller waited.** On 2026-10-06
|
|
// the verb took eighteen seconds on a mesh of four machines, so its caller read "still running" and
|
|
// had to ask `calls` for the answer to "is the mesh alright" — the one question that must answer at
|
|
// once, and the one a self-check and a rollout gate will ask every few minutes. So the serving
|
|
// controller composes it in the background — at its start, after anything that changes what it says
|
|
// (a machine's report, a build, a verb that acts), and every minute regardless — and the verb answers
|
|
// the last composition at once, saying when it was composed and how long that took. A caller who
|
|
// needs it newer than that reads the time and asks again; nothing is answered as current that is not.
|
|
|
|
// statusEvery is how often the summary is composed with nothing having nudged it; statusSettle how
|
|
// long a nudge waits for the next, so a push answered by four machines is composed once.
|
|
var (
|
|
statusEvery = time.Minute
|
|
statusSettle = 2 * time.Second
|
|
// statusComposeWithin bounds one composition, so a store that hangs cannot stop the summary for
|
|
// good; the attempt is said as failed, and the last summary stands with its age.
|
|
statusComposeWithin = 2 * time.Minute
|
|
)
|
|
|
|
// statusSummary is the last composed `status --json`, and when.
|
|
type statusSummary struct {
|
|
compose func(context.Context) ([]byte, error)
|
|
// every, settle and within are the clocks above, read once when it is made.
|
|
every, settle, within time.Duration
|
|
|
|
mu sync.Mutex
|
|
body []byte
|
|
composedAt time.Time
|
|
took time.Duration
|
|
failed string
|
|
failedAt time.Time
|
|
started time.Time
|
|
first chan struct{} // closed when the first attempt ends, either way
|
|
|
|
nudged chan struct{}
|
|
}
|
|
|
|
func newStatusSummary(compose func(context.Context) ([]byte, error)) *statusSummary {
|
|
return &statusSummary{compose: compose, started: time.Now(), first: make(chan struct{}),
|
|
nudged: make(chan struct{}, 1), every: statusEvery, settle: statusSettle, within: statusComposeWithin}
|
|
}
|
|
|
|
// statusFrom is the serving controller's summary; nil in any other process, where `status` is
|
|
// composed when asked, as at a shell.
|
|
var statusFrom *statusSummary
|
|
|
|
// nudge asks for a composition soon. Never blocks: one pending is as good as many.
|
|
func (s *statusSummary) nudge() {
|
|
if s == nil {
|
|
return
|
|
}
|
|
select {
|
|
case s.nudged <- struct{}{}:
|
|
default:
|
|
}
|
|
}
|
|
|
|
// keep composes until ctx ends: now, on a nudge once things settle, and every statusEvery.
|
|
func (s *statusSummary) keep(ctx context.Context) {
|
|
once := sync.Once{}
|
|
for {
|
|
s.composeOnce(ctx)
|
|
once.Do(func() { close(s.first) })
|
|
timer := time.NewTimer(s.every)
|
|
select {
|
|
case <-ctx.Done():
|
|
timer.Stop()
|
|
return
|
|
case <-timer.C:
|
|
case <-s.nudged:
|
|
timer.Stop()
|
|
// Let what else is arriving arrive, then compose once for all of it.
|
|
select {
|
|
case <-ctx.Done():
|
|
return
|
|
case <-time.After(s.settle):
|
|
}
|
|
select {
|
|
case <-s.nudged:
|
|
default:
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
func (s *statusSummary) composeOnce(ctx context.Context) {
|
|
start := time.Now()
|
|
asking, cancel := context.WithTimeout(ctx, s.within)
|
|
body, err := s.compose(asking)
|
|
cancel()
|
|
took := time.Since(start)
|
|
s.mu.Lock()
|
|
defer s.mu.Unlock()
|
|
if err != nil {
|
|
s.failed, s.failedAt = err.Error(), time.Now()
|
|
fmt.Printf("status could not be composed (after %s): %v — `status` answers the last summary, "+
|
|
"with its age\n", took.Round(time.Millisecond), err)
|
|
return
|
|
}
|
|
s.body, s.composedAt, s.took, s.failed = body, start, took, ""
|
|
}
|
|
|
|
// answer is what the `status` verb answers: the last summary at once, the same document `status
|
|
// --json` prints, with when it was composed. Before the first composition has ended it waits for it,
|
|
// but never past the caller's window; a controller that has none says so and why, rather than
|
|
// answering an empty mesh as a well one.
|
|
func (s *statusSummary) answer(ctx context.Context) (any, error) {
|
|
wait := time.NewTimer(link.AnswerWithin - time.Second)
|
|
defer wait.Stop()
|
|
select {
|
|
case <-s.first:
|
|
case <-wait.C:
|
|
case <-ctx.Done():
|
|
}
|
|
s.mu.Lock()
|
|
defer s.mu.Unlock()
|
|
if s.body == nil {
|
|
why := "its first composition has not finished"
|
|
if s.failed != "" {
|
|
why = "it could not be composed: " + s.failed
|
|
}
|
|
return nil, fmt.Errorf("this controller started %s ago and has no status to answer yet — %s. "+
|
|
"Ask again shortly", time.Since(s.started).Round(time.Second), why)
|
|
}
|
|
var parsed any
|
|
_ = json.Unmarshal(s.body, &parsed)
|
|
out := map[string]any{
|
|
"output": string(s.body), "ok": true, "answer": parsed,
|
|
"composed": s.composedAt.UTC().Format(time.RFC3339),
|
|
"age": time.Since(s.composedAt).Round(time.Second).String(),
|
|
"composedIn": s.took.Round(time.Millisecond).String(),
|
|
"note": "composed by the serving controller at its start, after each report, build or act, and " +
|
|
"every minute; answered at once from the last composition",
|
|
}
|
|
if s.failed != "" && s.failedAt.After(s.composedAt) {
|
|
out["lastAttemptFailed"] = fmt.Sprintf("%s: %s — this summary is the last that could be composed",
|
|
s.failedAt.UTC().Format(time.RFC3339), s.failed)
|
|
}
|
|
return out, nil
|
|
}
|
|
|
|
// composeStatus is `status --json`, composed in this process against its stores.
|
|
func composeStatus(open *stores) func(context.Context) ([]byte, error) {
|
|
return func(ctx context.Context) ([]byte, error) {
|
|
asked, err := theThreeQuestions(ctx, open)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
return statusAsJSON(asked)
|
|
}
|
|
}
|
|
|
|
// readingVerbs are the verbs that only read; after any other, what `status` says may have changed, so the summary is
|
|
// composed again; a verb that only reads leaves it alone, or a console polling `nodes` would keep the
|
|
// controller composing for ever.
|
|
var readingVerbs = map[string]bool{
|
|
"tools": true, "calls": true, "status": true, "nodes": true, "node": true, "modules": true,
|
|
"seats": true, "builds": true, "plan": true, "queue": true, "durations": true, "hand-acts": true,
|
|
"doctor": true, "conditions": true, "healers": true,
|
|
}
|
|
|
|
// nudgingListener is the enrolment, nudging the summary when a machine said something new.
|
|
type nudgingListener struct {
|
|
link.Enrolment
|
|
summary *statusSummary
|
|
// open is the serving controller's stores, for replacing a given value after a module's first
|
|
// good start (novox/hq ADR 0228).
|
|
open *stores
|
|
}
|
|
|
|
func (l nudgingListener) Heard(ctx context.Context, report link.Report) (bool, error) {
|
|
// A declaration refused as older than the one the machine holds is counted by its writer — the
|
|
// controller epoch it claimed (novox/hq to-be 45 §6, S13) — and so is what the machine's own count
|
|
// says it refused beyond the refusals heard.
|
|
now := time.Now()
|
|
if report.StaleRefusalOf() {
|
|
link.StaleRefusals.Refused(link.Refusal{Writer: link.WriterEpoch(report.Epoch), Epoch: report.Epoch,
|
|
Receiver: report.Node, At: now})
|
|
}
|
|
if report.Ordered() {
|
|
link.StaleRefusals.Lifetime(report.Node, report.RefusedOlder, now)
|
|
}
|
|
// What the machine's witnesses put back and stand by (novox/hq ADR 0236): read by the gate and its
|
|
// probe. Only from an account of the machine — not a word that a declaration was set aside, nor a rekey.
|
|
if report.Superseded == "" && report.Rekey == nil && report.Node != "" {
|
|
witnessed.heard(report.Node, report.Rollbacks, now)
|
|
}
|
|
news, err := l.Enrolment.Heard(ctx, report)
|
|
if news {
|
|
l.summary.nudge()
|
|
}
|
|
// What it says of its long-running resources (novox/hq ADR 0240): kept, and its modules' conditions
|
|
// raised or cleared from it. Only from an account of the machine the store took.
|
|
if err == nil && report.Health != nil && report.Superseded == "" && report.Rekey == nil && report.Node != "" &&
|
|
l.Enrolment.Inventory != nil {
|
|
if herr := stateHealth(ctx, l.Enrolment.Inventory, conditionsFrom, report.Node, *report.Health, now); herr != nil {
|
|
fmt.Fprintf(os.Stderr, "mesh-controller: could not keep what %s says of its resources' health: %v\n",
|
|
report.Node, herr)
|
|
}
|
|
}
|
|
if err == nil && l.open != nil && startedWell(report) {
|
|
// Off the report's path: replacing a given value sends the machine, and a report waits for
|
|
// nothing it caused (novox/hq ADR 0228).
|
|
go replaceGiven(context.WithoutCancel(ctx), l.open, report)
|
|
}
|
|
return news, err
|
|
}
|