Keep what each machine says of what it runs, raise it, and gate on it (hq ADR 0240, to-be 48 Phase A)
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.
This commit is contained in:
@@ -86,7 +86,9 @@ var WritersTable = []WriterRow{
|
||||
Others: "read", Subjects: []string{"mesh.node.*.declare"}, Writes: isController},
|
||||
{State: "a machine's applied state and its report", Writer: "the node-engine's apply queue",
|
||||
KeptIn: "the machine; the report on the bus", Others: "the reconcile and a delivery enqueue, never apply",
|
||||
Subjects: []string{"mesh.control.*.report"}, Writes: ownMachine},
|
||||
// And its health statement between reports (novox/hq ADR 0240): the same writer stating the same
|
||||
// machine, inside the grant it already had (`mesh.control.<its own>.>`).
|
||||
Subjects: []string{"mesh.control.*.report", "mesh.control.*.health"}, Writes: ownMachine},
|
||||
{State: "the controller lease", Writer: "the controller instance holding it", KeptIn: "key-value " + LeaseBucket,
|
||||
Others: "a candidate waits", Subjects: kvOf(LeaseBucket), Writes: isController},
|
||||
{State: "plans and their tiers", Writer: "controller (lease holder), compare-and-set on the plan's revision",
|
||||
|
||||
@@ -46,11 +46,14 @@ const (
|
||||
ScopeMesh = "mesh"
|
||||
// ScopeDelivery is a delivery mesh-delivery owns, by its id (novox/hq ADR 0239).
|
||||
ScopeDelivery = "delivery"
|
||||
// ScopeModule is a module on a machine, by `<module>.<machine>`: what it runs is not healthy there
|
||||
// (novox/hq ADR 0240).
|
||||
ScopeModule = "module"
|
||||
)
|
||||
|
||||
// Scopes is every scope, in the order a person reads them.
|
||||
var Scopes = []string{ScopeMachine, ScopePlan, ScopeCall, ScopeBuild, ScopeMerge, ScopeProvider,
|
||||
ScopeSeat, ScopeBus, ScopeCore, ScopeProbe, ScopeMesh, ScopeDelivery}
|
||||
ScopeSeat, ScopeBus, ScopeCore, ScopeProbe, ScopeMesh, ScopeDelivery, ScopeModule}
|
||||
|
||||
// Who resolves a condition.
|
||||
const (
|
||||
|
||||
@@ -0,0 +1,120 @@
|
||||
package inventory
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"time"
|
||||
|
||||
"github.com/jackc/pgx/v5"
|
||||
)
|
||||
|
||||
// What each machine says of its long-running resources (novox/hq ADR 0240, to-be 48 §4): the newest
|
||||
// statement per machine, and per module how many statements in a row said a resource of it was unhealthy.
|
||||
|
||||
// ResourceHealth is one long-running resource's state as the node-engine said it.
|
||||
type ResourceHealth struct {
|
||||
Module string `json:"module"`
|
||||
Resource string `json:"resource"`
|
||||
Kind string `json:"kind"`
|
||||
Target string `json:"target"`
|
||||
State string `json:"state"`
|
||||
Reason string `json:"reason,omitempty"`
|
||||
Since time.Time `json:"since"`
|
||||
Streak int `json:"streak,omitempty"`
|
||||
Restarts int `json:"restarts,omitempty"`
|
||||
}
|
||||
|
||||
// NodeHealth is a machine's newest statement, as kept.
|
||||
type NodeHealth struct {
|
||||
Node string
|
||||
Contract int
|
||||
// SaidAt is when the engine looked (the machine's clock); HeardAt when this controller heard it.
|
||||
SaidAt time.Time
|
||||
HeardAt time.Time
|
||||
Resources []ResourceHealth
|
||||
// Streaks is, per module, how many statements in a row said a resource of it was unhealthy.
|
||||
Streaks map[string]int
|
||||
}
|
||||
|
||||
// HealthOf is a machine's newest statement; false when its node-engine has never stated one.
|
||||
func (i *Inventory) HealthOf(ctx context.Context, nodeName string) (NodeHealth, bool, error) {
|
||||
all, err := i.healths(ctx, nodeName)
|
||||
if err != nil {
|
||||
return NodeHealth{}, false, err
|
||||
}
|
||||
h, ok := all[nodeName]
|
||||
return h, ok, nil
|
||||
}
|
||||
|
||||
// Healths is every machine's newest statement, by machine. A machine absent never stated one: its
|
||||
// node-engine is older than the judging, and its health is not known.
|
||||
func (i *Inventory) Healths(ctx context.Context) (map[string]NodeHealth, error) {
|
||||
return i.healths(ctx, "")
|
||||
}
|
||||
|
||||
func (i *Inventory) healths(ctx context.Context, only string) (map[string]NodeHealth, error) {
|
||||
rows, err := i.store.Pool().Query(ctx,
|
||||
`select n.name, h.contract, h.said_at, h.heard_at, h.resources, h.streaks
|
||||
from node_health h join node n on n.id = h.node
|
||||
where $1 = '' or n.name = $1`, only)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
defer rows.Close()
|
||||
out := map[string]NodeHealth{}
|
||||
for rows.Next() {
|
||||
var h NodeHealth
|
||||
var resources, streaks []byte
|
||||
if err := rows.Scan(&h.Node, &h.Contract, &h.SaidAt, &h.HeardAt, &resources, &streaks); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if err := json.Unmarshal(resources, &h.Resources); err != nil {
|
||||
return nil, fmt.Errorf("%s's health cannot be read: %w", h.Node, err)
|
||||
}
|
||||
if err := json.Unmarshal(streaks, &h.Streaks); err != nil {
|
||||
return nil, fmt.Errorf("%s's health cannot be read: %w", h.Node, err)
|
||||
}
|
||||
out[h.Node] = h
|
||||
}
|
||||
return out, rows.Err()
|
||||
}
|
||||
|
||||
// RecordHealth keeps a machine's statement in place of the one kept — unless the one kept is newer, by
|
||||
// when the engine looked: then nothing is written, and false says it was refused as older.
|
||||
func (i *Inventory) RecordHealth(ctx context.Context, h NodeHealth) (bool, error) {
|
||||
if h.Resources == nil {
|
||||
h.Resources = []ResourceHealth{}
|
||||
}
|
||||
if h.Streaks == nil {
|
||||
h.Streaks = map[string]int{}
|
||||
}
|
||||
resources, err := json.Marshal(h.Resources)
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
streaks, err := json.Marshal(h.Streaks)
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
heard := h.HeardAt
|
||||
if heard.IsZero() {
|
||||
heard = time.Now()
|
||||
}
|
||||
var node string
|
||||
err = i.store.Pool().QueryRow(ctx,
|
||||
`insert into node_health (node, contract, said_at, heard_at, resources, streaks)
|
||||
select id, $2, $3, $4, $5, $6 from node where name = $1
|
||||
on conflict (node) do update set contract = excluded.contract, said_at = excluded.said_at,
|
||||
heard_at = excluded.heard_at, resources = excluded.resources, streaks = excluded.streaks
|
||||
where node_health.said_at <= excluded.said_at
|
||||
returning node`, h.Node, h.Contract, h.SaidAt, heard, resources, streaks).Scan(&node)
|
||||
if errors.Is(err, pgx.ErrNoRows) {
|
||||
if _, nerr := i.NodeByName(ctx, h.Node); nerr != nil {
|
||||
return false, nerr
|
||||
}
|
||||
return false, nil
|
||||
}
|
||||
return err == nil, err
|
||||
}
|
||||
@@ -0,0 +1,26 @@
|
||||
-- A module says how it is healthy, and the node-engine judges it (novox/hq ADR 0240, to-be 48 Phase A).
|
||||
--
|
||||
-- Every machine's node-engine states the health of every long-running resource it runs for a module — a
|
||||
-- container that stays up, a process that stays up, a service stated running — in every report and as an
|
||||
-- event between reports. The controller keeps the newest statement per machine, here, so the release
|
||||
-- gate, `node show` and a controller started again all read the same word; and with it, per module, how
|
||||
-- many statements in a row said a resource of it was unhealthy — `module.<module>.<machine>.unhealthy` is
|
||||
-- raised on the second (ADR 0240 §4).
|
||||
--
|
||||
-- One row per machine, replaced, as node_report is: the question is the machine's state now. A machine
|
||||
-- whose node-engine is older than the judging has no row, and its health is not known — never healthy,
|
||||
-- never unhealthy.
|
||||
create table node_health (
|
||||
node uuid primary key references node(id) on delete cascade,
|
||||
-- The statement's version (the engine's liveness contract).
|
||||
contract int not null,
|
||||
-- When the node-engine looked, on the machine's clock: the order of its statements. An older one
|
||||
-- than this is refused.
|
||||
said_at timestamptz not null,
|
||||
-- When this controller heard it, on its own: what "since the send" is judged by.
|
||||
heard_at timestamptz not null default now(),
|
||||
-- [{module, resource, kind, target, state, reason, since, streak, restarts}], as the engine said them.
|
||||
resources jsonb not null default '[]',
|
||||
-- {module: statements in a row saying a resource of it is unhealthy}.
|
||||
streaks jsonb not null default '{}'
|
||||
);
|
||||
@@ -75,8 +75,16 @@ const (
|
||||
// ToolsAliveSubjects is every machine's node tools saying they are there (novox/hq to-be 45 §3,
|
||||
// S11): core NATS like the host's, for the same reason.
|
||||
ToolsAliveSubjects = "mesh.control.*.tools-alive"
|
||||
|
||||
// HealthSubjects is every machine's health statement between its reports (novox/hq ADR 0240): core
|
||||
// NATS like the heartbeat, because a statement lost is said again within a minute while anything is
|
||||
// not healthy, and the next report carries it whatever happens.
|
||||
HealthSubjects = "mesh.control.*.health"
|
||||
)
|
||||
|
||||
// HealthSubject is one machine's health statement.
|
||||
func HealthSubject(node string) string { return "mesh.control." + node + ".health" }
|
||||
|
||||
// ReportSubject is where one node says what it did. On the CONTROL stream, because it is the
|
||||
// message the store-window guarantee is about (ADR 0083).
|
||||
func ReportSubject(node string) string { return "mesh.control." + node + ".report" }
|
||||
|
||||
@@ -35,6 +35,10 @@ var Contracts = map[string]Contract{
|
||||
KindHeartbeat: {Unordered: "a word that the machine is there: the newest heard is the newest said, and one " +
|
||||
"lost is the next one"},
|
||||
KindToolsHeartbeat: {Unordered: "a word that the node tools are there, as a machine's heartbeat"},
|
||||
KindHealth: {Ordered: "by the time the node-engine looked (`at`, on the machine's clock): a statement older " +
|
||||
"than the one kept for that machine is refused, so an event arriving after a newer report does not undo it " +
|
||||
"(novox/hq ADR 0240)",
|
||||
Tests: []string{"TestAnOlderHealthStatementIsRefused"}},
|
||||
KindEnrolment: {Unordered: "a request answered once, under a token spent once: a second presentation is " +
|
||||
"refused by the token, not by an order (issue 083)",
|
||||
Tests: []string{"TestAnEnrolmentMetByAHeldTokenIsAskedToTryAgain"}},
|
||||
|
||||
@@ -16,7 +16,7 @@ import (
|
||||
func TestEveryConsumedKindHasAContract(t *testing.T) {
|
||||
// What the controller can be handed, from the subjects it is granted and the ones it derives.
|
||||
subjects := []string{EnrolSubject, BuiltSubject, ReportSubject("anchor"), AliveSubject("anchor"),
|
||||
ToolsAliveSubject("anchor"), "mesh.mod.postgres.event.provisioner.failing"}
|
||||
ToolsAliveSubject("anchor"), HealthSubject("anchor"), "mesh.mod.postgres.event.provisioner.failing"}
|
||||
subjects = append(subjects, broker.ControllerFollows...)
|
||||
kinds := map[string]bool{}
|
||||
for _, s := range subjects {
|
||||
|
||||
@@ -0,0 +1,56 @@
|
||||
package link
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"strings"
|
||||
)
|
||||
|
||||
// A machine's health statement (novox/hq ADR 0240, to-be 48 §4).
|
||||
//
|
||||
// **The node-engine owns every verdict; the controller keeps the last word per machine.** A statement
|
||||
// arrives in every report and as an event between reports — on each change, and again every minute while
|
||||
// a resource is not healthy. What the controller does with it — keep it, refuse an older one, raise
|
||||
// `module.<module>.<machine>.unhealthy` on the second statement that says so — is the Healths it is given.
|
||||
|
||||
// Healths keeps what machines state of their long-running resources.
|
||||
type Healths interface {
|
||||
// Stated keeps one machine's statement, from a report or the event. An older statement than the one
|
||||
// kept is refused there, by its time on the machine.
|
||||
Stated(ctx context.Context, node string, h Health) error
|
||||
}
|
||||
|
||||
// Hears says where the machines' health statements are kept. Their subscription is the heartbeats':
|
||||
// core, and always made, so nothing is asked of the bus here.
|
||||
func (s *Server) Hears(h Healths) { s.healths = h }
|
||||
|
||||
// healthSaid acts on one health event. Core NATS, so there is nothing to hold: a statement that could not
|
||||
// be kept is said in the log, and the next one — a minute away while anything is not healthy — is kept.
|
||||
func (s *Server) healthSaid(ctx context.Context, m Control) {
|
||||
defer func() { _ = m.Took() }()
|
||||
var said HealthSaid
|
||||
if err := json.Unmarshal(m.Body(), &said); err != nil || said.Health.Contract == 0 {
|
||||
return
|
||||
}
|
||||
// The machine is the one in the subject the bus let it publish on, never the body's.
|
||||
node, ok := nodeOfHealth(m.Subject())
|
||||
if !ok || s.healths == nil {
|
||||
return
|
||||
}
|
||||
if err := s.healths.Stated(ctx, node, said.Health); err != nil {
|
||||
s.log.Printf("could not keep what %s says of its resources' health: %v", node, err)
|
||||
}
|
||||
}
|
||||
|
||||
// nodeOfHealth is the machine a health statement names in its subject.
|
||||
func nodeOfHealth(subject string) (string, bool) {
|
||||
rest, ok := strings.CutPrefix(subject, "mesh.control.")
|
||||
if !ok {
|
||||
return "", false
|
||||
}
|
||||
node, kind, ok := strings.Cut(rest, ".")
|
||||
if !ok || kind != "health" || node == "" || strings.Contains(node, ".") {
|
||||
return "", false
|
||||
}
|
||||
return node, true
|
||||
}
|
||||
@@ -0,0 +1,37 @@
|
||||
package link
|
||||
|
||||
import (
|
||||
"context"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
// A health event is kept as the machine its subject names says it, whatever its body claims, and taken
|
||||
// (novox/hq ADR 0240): core NATS, nothing to hold.
|
||||
|
||||
type keptHealths struct{ by map[string]Health }
|
||||
|
||||
func (k *keptHealths) Stated(_ context.Context, node string, h Health) error {
|
||||
k.by[node] = h
|
||||
return nil
|
||||
}
|
||||
|
||||
func TestAHealthEventIsKeptAsTheMachineInItsSubject(t *testing.T) {
|
||||
s, in := serving()
|
||||
kept := &keptHealths{by: map[string]Health{}}
|
||||
s.Hears(kept)
|
||||
to := &settled{}
|
||||
m := in.sends(t, to, KindHealth, HealthSaid{Node: "laptop", Health: Health{Contract: LivenessContract, At: time.Now(),
|
||||
Resources: []ResourceHealth{{Module: "letta", Resource: "letta.server", State: StateUnhealthy}}}}).(*fakeControl)
|
||||
m.subject = HealthSubject("anchor")
|
||||
s.act(t.Context(), m)
|
||||
if !to.acked {
|
||||
t.Fatal("a health event was not taken")
|
||||
}
|
||||
if _, lied := kept.by["laptop"]; lied || len(kept.by["anchor"].Resources) != 1 {
|
||||
t.Fatalf("kept %+v: the machine is the subject's, never the body's", kept.by)
|
||||
}
|
||||
if kind, ok := kindOfSubject(HealthSubject("anchor")); !ok || kind != KindHealth {
|
||||
t.Fatalf("%s is not read as a health statement", HealthSubject("anchor"))
|
||||
}
|
||||
}
|
||||
@@ -262,6 +262,13 @@ type Report struct {
|
||||
// `witness` and `not-reversible` are sent only to a machine whose report carries it.
|
||||
Witness int `json:"witness,omitempty"`
|
||||
|
||||
// Health is the machine's word on every long-running resource it runs for a module (novox/hq ADR
|
||||
// 0240, to-be 48 §4; mesh-host's internal/liveness): its state, since when, its failing streak and
|
||||
// the restarts its node-engine counted. **Absent from an engine older than the judging**, which is
|
||||
// read as "not known" — never as healthy, never as a reason to raise anything; present with no
|
||||
// resources from a machine that runs nothing long-lived.
|
||||
Health *Health `json:"health,omitempty"`
|
||||
|
||||
// Rekey is a node taking a found tunnel's key as its overlay key after enrolment (novox/hq
|
||||
// ADR 0105). A report carrying one is not an account of the machine: it moves the node's
|
||||
// overlay key and tunnel and nothing else.
|
||||
@@ -404,3 +411,45 @@ func EnrolProof(secret string, public []byte, overlay, sealing, serving string)
|
||||
return []byte("novox-mesh-enrol\x00" + secret + "\x00" + base64.StdEncoding.EncodeToString(public) +
|
||||
"\x00" + overlay + "\x00" + sealing + "\x00" + serving)
|
||||
}
|
||||
|
||||
// LivenessContract is the version of the health statement this controller reads (ADR 0240 Phase A).
|
||||
const LivenessContract = 1
|
||||
|
||||
// Health is one statement of a machine's long-running resources (to-be 48 §4): in every report, as the
|
||||
// event HealthSubject between reports on each change, and again every minute while one is not healthy.
|
||||
// The node-engine's own (mesh-host internal/link Health); a test on each side holds the field names.
|
||||
type Health struct {
|
||||
Contract int `json:"contract"`
|
||||
// At is when the engine looked, on the machine's clock: the order of its statements.
|
||||
At time.Time `json:"at"`
|
||||
Resources []ResourceHealth `json:"resources"`
|
||||
}
|
||||
|
||||
// The states a resource is said in (ADR 0240 §4).
|
||||
const (
|
||||
StateHealthy = "healthy"
|
||||
StateUnhealthy = "unhealthy"
|
||||
StateStarting = "starting"
|
||||
StateHeld = "held"
|
||||
StateUnknown = "unknown"
|
||||
)
|
||||
|
||||
// ResourceHealth is one long-running resource's state.
|
||||
type ResourceHealth struct {
|
||||
Module string `json:"module"`
|
||||
Resource string `json:"resource"`
|
||||
Kind string `json:"kind"`
|
||||
Target string `json:"target"`
|
||||
State string `json:"state"`
|
||||
Reason string `json:"reason,omitempty"`
|
||||
Since time.Time `json:"since"`
|
||||
Streak int `json:"streak,omitempty"`
|
||||
Restarts int `json:"restarts,omitempty"`
|
||||
}
|
||||
|
||||
// HealthSaid is the health event's body: the machine and its statement. The machine is read from the
|
||||
// subject the bus let it publish on, never from here.
|
||||
type HealthSaid struct {
|
||||
Node string `json:"node"`
|
||||
Health Health `json:"health"`
|
||||
}
|
||||
|
||||
@@ -25,6 +25,13 @@ func TestTheWireFormatIsExactlyTheseFieldNames(t *testing.T) {
|
||||
[]string{"protocol", "address", "port", "by", "published", "container-port"}},
|
||||
{EnrolRequest{Node: "n", Secret: "s", PublicKey: []byte("k")},
|
||||
[]string{"node", "secret", "public_key"}},
|
||||
// novox/hq ADR 0240: every long-running resource's health, in every report and in its own event.
|
||||
{Report{Node: "n", Health: &Health{Contract: LivenessContract}}, []string{"node", "health"}},
|
||||
{Health{Contract: LivenessContract, Resources: []ResourceHealth{}}, []string{"contract", "at", "resources"}},
|
||||
{ResourceHealth{Module: "m", Resource: "m.r", Kind: "container", Target: "t", State: StateUnhealthy,
|
||||
Reason: "restarting", Streak: 2, Restarts: 3},
|
||||
[]string{"module", "resource", "kind", "target", "state", "reason", "since", "streak", "restarts"}},
|
||||
{HealthSaid{Node: "n"}, []string{"node", "health"}},
|
||||
} {
|
||||
raw, err := json.Marshal(c.value)
|
||||
if err != nil {
|
||||
|
||||
@@ -30,8 +30,10 @@ const (
|
||||
KindHeartbeat = "heartbeat"
|
||||
// KindToolsHeartbeat is a machine's node tools saying they are there (novox/hq to-be 45 S11).
|
||||
KindToolsHeartbeat = "tools-heartbeat"
|
||||
KindBuilt = "built"
|
||||
KindModuleMoved = "module-moved"
|
||||
// KindHealth is a machine's statement of its long-running resources' health (novox/hq ADR 0240).
|
||||
KindHealth = "health"
|
||||
KindBuilt = "built"
|
||||
KindModuleMoved = "module-moved"
|
||||
// KindSourceMoved is the forge announcing a merge: a source moved, and what it produces is
|
||||
// built without anybody telling the mesh (novox/hq 04-ISSUES/131).
|
||||
KindSourceMoved = "source-moved"
|
||||
|
||||
@@ -108,6 +108,13 @@ func (n *natsInbound) Receive(ctx context.Context, act func(context.Context, Con
|
||||
return fmt.Errorf("subscribing to the node tools' heartbeats: %w", err)
|
||||
}
|
||||
defer func() { _ = toolsAlive.Unsubscribe() }()
|
||||
// And what each machine says of its long-running resources between reports (novox/hq ADR 0240): core
|
||||
// too, and said again while it matters.
|
||||
health, err := conn.ChanSubscribe(HealthSubjects, beats)
|
||||
if err != nil {
|
||||
return fmt.Errorf("subscribing to the machines' health: %w", err)
|
||||
}
|
||||
defer func() { _ = health.Unsubscribe() }()
|
||||
|
||||
// The events the controller follows, when something is listening for them.
|
||||
var events chan *nats.Msg
|
||||
@@ -239,6 +246,8 @@ func kindOfSubject(subject string) (string, bool) {
|
||||
return KindHeartbeat, true
|
||||
case "tools-alive":
|
||||
return KindToolsHeartbeat, true
|
||||
case "health":
|
||||
return KindHealth, true
|
||||
}
|
||||
}
|
||||
switch subject {
|
||||
|
||||
@@ -94,6 +94,8 @@ type Server struct {
|
||||
retirements Retirements
|
||||
// checker asks for a pull request's merge check (novox/hq to-be 45 §9).
|
||||
checker Checker
|
||||
// healths keeps what machines say of their long-running resources (novox/hq ADR 0240).
|
||||
healths Healths
|
||||
|
||||
log *log.Logger
|
||||
// giveUp is how long one message is held for the store; zero means GiveUpAfter.
|
||||
@@ -194,6 +196,8 @@ func (s *Server) act(ctx context.Context, m Control) {
|
||||
s.heartbeat(m)
|
||||
case KindToolsHeartbeat:
|
||||
s.toolsHeartbeat(m)
|
||||
case KindHealth:
|
||||
s.healthSaid(ctx, m)
|
||||
case KindBuilt:
|
||||
s.wasBuilt(ctx, m)
|
||||
case KindModuleMoved:
|
||||
|
||||
Reference in New Issue
Block a user