Liveness alone could not see a web application whose port was open and whose program ran while every request hung for eleven hours (issue 145). A resource now carries the `health` its module declared: the engine makes http and tcp looks itself from the machine to the endpoint's published port, reads a unit's readiness from the show it already makes, hands an exec command or the image's own check to the runtime as the container's check with the declared timing and reads its state from the inspect it already makes, and asks a module's tool on its own node tools. Starting until the check passed, unhealthy once its looks after the grace fail the declared number of times; never more looks than the measured budget; nothing restarted. The statement says contract 2, which tells the controller this engine may be sent the field.
557 lines
19 KiB
Go
557 lines
19 KiB
Go
package link
|
|
|
|
import (
|
|
"context"
|
|
"crypto/sha256"
|
|
"encoding/hex"
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"sync"
|
|
"time"
|
|
)
|
|
|
|
// Queue is the node-engine's one apply queue (novox/hq to-be 45 §6, ADR 0227 rule 1).
|
|
//
|
|
// **One worker applies; everything else asks it to.** A declaration delivered over the link and the
|
|
// five-minute reconcile are two reasons to *enqueue* the same act, never two paths that apply. Before
|
|
// this they were two, each taking an apply lock, and every order between them was a fault met on the
|
|
// live mesh: the reconcile read what was kept before a delivery and applied it over the newer one
|
|
// (novox/hq issues 257, 261); the reconcile went first and its report about the older declaration
|
|
// reached the mesh after the delivery's (267). Each fix closed one order and left the class open.
|
|
//
|
|
// The queue holds at most one pending request, coalesced: whatever has been delivered and not yet
|
|
// taken, and whether a reconcile is due. When the worker starts it takes the **newest declaration held
|
|
// at that moment**, sets the rest aside — each reported as set aside, then settled unapplied — applies
|
|
// it once, and makes **one report**, naming the declaration's order and its own report sequence. A
|
|
// reconcile due while a delivery waits is that delivery's apply: applying the newest thing the mesh
|
|
// said is what reconciling is. Only with nothing delivered does a reconcile apply what was kept, read
|
|
// once its turn has come.
|
|
//
|
|
// **Reports leave in the order they are made**, because the worker makes and says them, one at a time.
|
|
// A report made while the link is down waits: an apply's in the unsaid store (issue 264), a
|
|
// reconcile's in one slot replaced by anything newer, and both are said, in that order, when the next
|
|
// link opens — before anything newly delivered is applied.
|
|
type Queue struct {
|
|
// Membership is whose declarations these are: their signer, and the node the reports are from.
|
|
Membership Membership
|
|
// Apply applies a declaration the mesh signed.
|
|
Apply Applier
|
|
// Reconcile holds the machine to what it was last told, read once it is this act's turn. False
|
|
// when there is nothing to hold it to, or the reconcile could not run; it says why itself.
|
|
Reconcile Reconcile
|
|
// News is whether a reconcile's report is worth saying unasked (novox/hq ADR 0100, ADR 0140). Nil
|
|
// is never: a reconcile is otherwise silent.
|
|
News func(Report) bool
|
|
// Heard is told every account of the machine the mesh has taken — an apply's or a reconcile's,
|
|
// not a refusal or a declaration set aside, which say nothing about the machine. Nil is allowed.
|
|
Heard func(Report)
|
|
// Unsaid keeps an apply's report until the mesh has taken it (issue 264). Nil keeps nothing.
|
|
Unsaid Unsaid
|
|
// Numbers keeps the report sequence and the refusals counted, across restarts and self-updates.
|
|
// Nil counts in memory, which is only right where nothing outlives the process — a test.
|
|
Numbers Numbers
|
|
Say Announce
|
|
Timeout time.Duration
|
|
|
|
once sync.Once
|
|
wake chan struct{}
|
|
|
|
mu sync.Mutex
|
|
idle *sync.Cond
|
|
// delivering is a delivered declaration in hand: its report is owed on the link it came over, so
|
|
// that link is not let go until it is said (detach). A reconcile in hand holds no link — its report
|
|
// waits for the next one if this one goes.
|
|
delivering bool
|
|
waiting []Declaration // delivered and not yet taken, in the order they arrived
|
|
due bool // a reconcile asked for
|
|
// asked is the mesh asking this node to say what it last applied (the `report` verb, to-be 45 §6):
|
|
// the reconcile it enqueues says its report whether or not it is news.
|
|
asked bool
|
|
bus Bus // the link open now; nil while there is none
|
|
unasked *Report // a reconcile's report not yet said, made while there was no link
|
|
|
|
// saying is held while a report is made and said, so the order they leave in is the order they
|
|
// are made in — whoever says them, the worker or a link opening.
|
|
saying sync.Mutex
|
|
counted bool
|
|
numbers numbered
|
|
}
|
|
|
|
// Reconcile is the reconcile's act, run by the queue's worker. See Queue.Reconcile.
|
|
type Reconcile func(ctx context.Context) (Report, bool)
|
|
|
|
// Numbers is where the queue keeps what it counts, so a successor goes on from where it stopped.
|
|
type Numbers interface {
|
|
// Read is what was kept: zero for a node that never reported, an error when it cannot be read —
|
|
// never zero for "could not tell".
|
|
Read() (reportSequence, refusedOlder int64, err error)
|
|
// Save keeps them, on disk when it returns.
|
|
Save(reportSequence, refusedOlder int64) error
|
|
}
|
|
|
|
type numbered struct{ sequence, refused int64 }
|
|
|
|
func (q *Queue) init() {
|
|
q.once.Do(func() {
|
|
q.wake = make(chan struct{}, 1)
|
|
q.idle = sync.NewCond(&q.mu)
|
|
if q.Say == nil {
|
|
q.Say = func(string) {}
|
|
}
|
|
})
|
|
}
|
|
|
|
// Deliver enqueues what the link delivered, in the order it arrived. It never waits for an apply.
|
|
func (q *Queue) Deliver(arrived ...Declaration) {
|
|
q.init()
|
|
q.mu.Lock()
|
|
q.waiting = append(q.waiting, arrived...)
|
|
q.mu.Unlock()
|
|
q.signal()
|
|
}
|
|
|
|
// ReconcileDue enqueues a reconcile. One asked for while another is waiting is the same one.
|
|
func (q *Queue) ReconcileDue() {
|
|
q.init()
|
|
q.mu.Lock()
|
|
q.due = true
|
|
q.mu.Unlock()
|
|
q.signal()
|
|
}
|
|
|
|
// ReportAsked is the mesh asking this node to say again what it last applied (novox/hq to-be 45 §6, the
|
|
// `report` verb its healer H1 asks when a send went unreported). **An enqueue, like everything else**: a
|
|
// reconcile, whose report — the account of the declaration this node keeps, the last it applied — is
|
|
// said whether or not it is news. A delivery waiting meanwhile is applied instead and reported, which
|
|
// answers the question better. Asked again before the worker takes it is the same question.
|
|
func (q *Queue) ReportAsked() {
|
|
q.init()
|
|
q.mu.Lock()
|
|
q.due, q.asked = true, true
|
|
q.mu.Unlock()
|
|
q.signal()
|
|
}
|
|
|
|
func (q *Queue) signal() {
|
|
select {
|
|
case q.wake <- struct{}{}:
|
|
default:
|
|
// One is already waiting to be read, and the worker takes everything pending when it does.
|
|
}
|
|
}
|
|
|
|
// Run is the worker. It returns once the context ends and the act in hand is finished and said:
|
|
// asked to stop — the host standing aside for its successor — nothing further is applied.
|
|
func (q *Queue) Run(ctx context.Context) {
|
|
q.init()
|
|
for {
|
|
for q.step(ctx) {
|
|
}
|
|
select {
|
|
case <-ctx.Done():
|
|
return
|
|
case <-q.wake:
|
|
}
|
|
}
|
|
}
|
|
|
|
// step takes what is pending and does it, once. False when there was nothing to do, or the queue
|
|
// was asked to stop.
|
|
func (q *Queue) step(ctx context.Context) bool {
|
|
q.mu.Lock()
|
|
if ctx.Err() != nil || (len(q.waiting) == 0 && !q.due) {
|
|
q.mu.Unlock()
|
|
return false
|
|
}
|
|
batch, due, asked := q.waiting, q.due, q.asked
|
|
q.waiting, q.due, q.asked = nil, false, false
|
|
q.delivering = len(batch) > 0
|
|
q.mu.Unlock()
|
|
defer func() {
|
|
q.mu.Lock()
|
|
q.delivering = false
|
|
q.idle.Broadcast()
|
|
q.mu.Unlock()
|
|
}()
|
|
|
|
if len(batch) == 0 {
|
|
q.reconcile(ctx, asked)
|
|
return true
|
|
}
|
|
|
|
latest, aside := pick(batch)
|
|
for _, old := range aside {
|
|
if d := declaredIn(old.Body()); d != "" && d == declaredIn(latest.Body()) {
|
|
// The same declaration again — delivered twice while the worker was busy. Nothing is set
|
|
// aside: it is about to be applied.
|
|
_ = old.Handled()
|
|
continue
|
|
}
|
|
q.Say(fmt.Sprintf("set aside declaration %s (%s): a newer one arrived with it",
|
|
short(declaredIn(old.Body())), orderOf(old.Body()).Words()))
|
|
q.tell(ctx, Report{Declared: declaredIn(old.Body()), Order: orderOf(old.Body()),
|
|
Superseded: declaredIn(latest.Body())}, setAside, old)
|
|
}
|
|
|
|
report := handleBody(ctx, q.Membership, latest.Body(), q.Apply)
|
|
switch {
|
|
case report.OlderThan != nil:
|
|
// **Refused, counted and said** (rule 2): one line here, the count on this and every later
|
|
// report, and the report itself naming what was refused and what is held.
|
|
total := q.countRefusal()
|
|
q.Say(fmt.Sprintf("refused declaration %s (%s): this node applied %s, from a newer lease "+
|
|
"holder — %d refused as older so far", short(report.Declared), report.Order.Words(),
|
|
report.OlderThan.Words(), total))
|
|
case report.Refused != "":
|
|
q.Say("refused a declaration: " + report.Refused)
|
|
case len(report.Failed) > 0:
|
|
q.Say(fmt.Sprintf("applied %d and failed: %v%s",
|
|
len(report.Applied), report.Failed, heldNote(report.Held)))
|
|
default:
|
|
q.Say(fmt.Sprintf("applied %d resource(s)%s", len(report.Applied), heldNote(report.Held)))
|
|
}
|
|
if due && report.Refused != "" {
|
|
// Nothing was applied, so the reconcile this stood in for has not happened. It runs next.
|
|
q.mu.Lock()
|
|
q.due = true
|
|
q.mu.Unlock()
|
|
}
|
|
// A refusal as older is not kept to be said again: what the mesh waits for from this node is the
|
|
// account of the last apply, and that is still what is kept.
|
|
kind := applied
|
|
if report.OlderThan != nil {
|
|
kind = refusedOlder
|
|
}
|
|
q.tell(ctx, report, kind, latest)
|
|
return true
|
|
}
|
|
|
|
func (q *Queue) reconcile(ctx context.Context, asked bool) {
|
|
if q.Reconcile == nil {
|
|
if asked {
|
|
q.Say("asked to say what this node last applied, and it holds nothing to reconcile with")
|
|
}
|
|
return
|
|
}
|
|
report, ok := q.Reconcile(ctx)
|
|
if !ok {
|
|
if asked {
|
|
q.Say("asked to say what this node last applied, and the reconcile could not say it")
|
|
}
|
|
return
|
|
}
|
|
if !asked && (q.News == nil || !q.News(report)) {
|
|
return
|
|
}
|
|
if asked {
|
|
q.Say(fmt.Sprintf("asked by the mesh, saying again what this node last applied: declaration %s (%s)",
|
|
short(report.Declared), report.Order.Words()))
|
|
}
|
|
q.tell(ctx, report, unasked)
|
|
}
|
|
|
|
// What a report is, for what tell does with it besides saying it.
|
|
type kindOf int
|
|
|
|
const (
|
|
// applied is an apply's account: kept until said (issue 264), and fresher than any reconcile's
|
|
// report still waiting.
|
|
applied kindOf = iota
|
|
// setAside is a declaration a newer one took the place of, unapplied.
|
|
setAside
|
|
// refusedOlder is a declaration refused as older than what was applied. Not kept: what the mesh
|
|
// waits for from this node is still the account of the last apply.
|
|
refusedOlder
|
|
// unasked is a reconcile's report. Not said — no link, or the bus refused it — the newest waits
|
|
// for the next link, and anything the worker says before then replaces it.
|
|
unasked
|
|
)
|
|
|
|
// tell makes a report this node's — its node, its report sequence, the refusals counted — and says it
|
|
// on the link open now. Kept first, when it is an apply's, so a report lost between here and the bus
|
|
// is said by whichever host links next. Each declaration it settles is settled after it is said,
|
|
// either way. True when the broker took it.
|
|
func (q *Queue) tell(ctx context.Context, r Report, kind kindOf, settle ...Declaration) bool {
|
|
q.saying.Lock()
|
|
defer q.saying.Unlock()
|
|
|
|
r = q.stamp(r)
|
|
keep := kind == applied
|
|
if keep {
|
|
// What this says supersedes any reconcile's report still waiting to be said.
|
|
q.unasked = nil
|
|
if q.Unsaid != nil {
|
|
if err := q.Unsaid.Keep(r); err != nil {
|
|
q.Say("cannot keep this apply's report until it is said: " + err.Error())
|
|
}
|
|
}
|
|
}
|
|
q.mu.Lock()
|
|
bus := q.bus
|
|
q.mu.Unlock()
|
|
|
|
published := false
|
|
if bus != nil {
|
|
// **Said even when the link is being let go** (issue 264): not cancelled with it, still
|
|
// bounded by the timeout, and the link is not closed until this returns (detach).
|
|
published = publishReport(context.WithoutCancel(ctx), bus, q.Membership, r, q.Say, q.Timeout)
|
|
}
|
|
if !published && kind == unasked {
|
|
q.unasked = &r
|
|
}
|
|
if published {
|
|
if keep && q.Unsaid != nil {
|
|
if err := q.Unsaid.Said(r.Declared); err != nil {
|
|
q.Say("said this apply's report, and cannot forget it: " + err.Error())
|
|
}
|
|
}
|
|
if q.Heard != nil && (kind == applied || kind == unasked) && r.Refused == "" {
|
|
q.Heard(r)
|
|
}
|
|
}
|
|
// Settled after the report is published. A node that dies between applying and reporting
|
|
// leaves the declaration with the mesh and applies it again on return, which is safe because
|
|
// applying is reconciliation — it converges rather than repeating.
|
|
for _, d := range settle {
|
|
_ = d.Handled()
|
|
}
|
|
return published
|
|
}
|
|
|
|
// stamp gives a report its node, the next report sequence and the refusals counted. Called holding
|
|
// saying.
|
|
func (q *Queue) stamp(r Report) Report {
|
|
q.load()
|
|
q.numbers.sequence++
|
|
q.keepNumbers()
|
|
r.Node = q.Membership.Node
|
|
r.ReportSequence = q.numbers.sequence
|
|
r.RefusedOlder = q.numbers.refused
|
|
return r
|
|
}
|
|
|
|
func (q *Queue) countRefusal() int64 {
|
|
q.saying.Lock()
|
|
defer q.saying.Unlock()
|
|
q.load()
|
|
q.numbers.refused++
|
|
q.keepNumbers()
|
|
return q.numbers.refused
|
|
}
|
|
|
|
// load reads what was counted, once. **Unreadable is said, and the sequence goes on from above
|
|
// anything it can have reached** — the time in milliseconds — rather than from one, which would make
|
|
// every report this host says next older than what it already said, and refused.
|
|
func (q *Queue) load() {
|
|
if q.counted {
|
|
return
|
|
}
|
|
q.counted = true
|
|
if q.Numbers == nil {
|
|
return
|
|
}
|
|
sequence, refused, err := q.Numbers.Read()
|
|
if err != nil {
|
|
floor := time.Now().UnixMilli()
|
|
q.Say(fmt.Sprintf("cannot read this node's report sequence: %v — going on from %d, above any "+
|
|
"it can have reached, and counting refusals again from none", err, floor))
|
|
q.numbers = numbered{sequence: floor}
|
|
return
|
|
}
|
|
q.numbers = numbered{sequence: sequence, refused: refused}
|
|
}
|
|
|
|
func (q *Queue) keepNumbers() {
|
|
if q.Numbers == nil {
|
|
return
|
|
}
|
|
if err := q.Numbers.Save(q.numbers.sequence, q.numbers.refused); err != nil {
|
|
// Said, and not fatal: the report is still worth saying. What is lost is that a successor
|
|
// could number one again.
|
|
q.Say("cannot keep this node's report sequence: " + err.Error())
|
|
}
|
|
}
|
|
|
|
// attach is a link opening. **What the last apply did, if the mesh never heard it**, is said first,
|
|
// exactly as it was made; then a reconcile's report that waited; and only then may the worker say
|
|
// anything on it — so a report said again can never land after, and read as newer than, a later one.
|
|
func (q *Queue) attach(ctx context.Context, bus Bus) {
|
|
q.init()
|
|
q.saying.Lock()
|
|
defer q.saying.Unlock()
|
|
reporting := context.WithoutCancel(ctx)
|
|
q.load()
|
|
|
|
said := true
|
|
if q.Unsaid != nil {
|
|
switch kept, ok, err := q.Unsaid.Pending(); {
|
|
case err != nil:
|
|
q.Say("cannot read the report this node kept unsaid: " + err.Error())
|
|
case ok:
|
|
q.Say("saying again what the last apply did: its report never reached the mesh")
|
|
if said = publishReport(reporting, bus, q.Membership, kept, q.Say, q.Timeout); said {
|
|
if err := q.Unsaid.Said(kept.Declared); err != nil {
|
|
q.Say("said the kept report, and cannot forget it: " + err.Error())
|
|
}
|
|
if q.Heard != nil && kept.Refused == "" {
|
|
q.Heard(kept)
|
|
}
|
|
}
|
|
}
|
|
}
|
|
if said && q.unasked != nil {
|
|
if publishReport(reporting, bus, q.Membership, *q.unasked, q.Say, q.Timeout) {
|
|
if q.Heard != nil {
|
|
q.Heard(*q.unasked)
|
|
}
|
|
q.unasked = nil
|
|
}
|
|
}
|
|
|
|
q.mu.Lock()
|
|
q.bus = bus
|
|
q.mu.Unlock()
|
|
}
|
|
|
|
// detach is the link closing: it waits for a delivery in hand to be applied and said — an apply that
|
|
// stood aside for its successor is reported on this link before it is let go (issue 264) — and then
|
|
// nothing more is said on it. A reconcile in hand is not waited for: a link that was dropped or roused
|
|
// is opened again at once, and the reconcile's report, if it is news, waits for it.
|
|
func (q *Queue) detach(bus Bus) {
|
|
q.init()
|
|
q.mu.Lock()
|
|
for q.delivering {
|
|
q.idle.Wait()
|
|
}
|
|
q.mu.Unlock()
|
|
|
|
q.saying.Lock()
|
|
defer q.saying.Unlock()
|
|
q.mu.Lock()
|
|
if q.bus == bus {
|
|
q.bus = nil
|
|
}
|
|
q.mu.Unlock()
|
|
}
|
|
|
|
// pick is the newest of what was delivered and the rest, in arrival order, to set aside: by order
|
|
// when the declarations claim one, by arrival when they do not (Order.Supersedes).
|
|
func pick(batch []Declaration) (Declaration, []Declaration) {
|
|
latest := batch[0]
|
|
var aside []Declaration
|
|
for _, next := range batch[1:] {
|
|
if orderOf(next.Body()).Supersedes(orderOf(latest.Body())) {
|
|
aside = append(aside, latest)
|
|
latest = next
|
|
continue
|
|
}
|
|
aside = append(aside, next)
|
|
}
|
|
return latest, aside
|
|
}
|
|
|
|
// orderOf is the order a signed declaration claims, or none when it claims none or cannot be read.
|
|
// Read from the envelope alone; the signature is verified when the winner is applied, and a forged
|
|
// message that lied about its order would only set aside real ones — which are reported as set aside,
|
|
// and the next push sends the current one again.
|
|
func orderOf(body []byte) Order {
|
|
var signed Signed
|
|
if err := json.Unmarshal(body, &signed); err != nil {
|
|
return Order{}
|
|
}
|
|
var o Order
|
|
if err := json.Unmarshal(signed.Declaration, &o); err != nil || o.Epoch < 0 || o.Sequence < 0 {
|
|
return Order{}
|
|
}
|
|
return o
|
|
}
|
|
|
|
// declaredIn names a signed declaration as the mesh does — the digest of the declaration's bytes, the
|
|
// same the apply's report carries — for a report about one that was not applied. Empty if the message
|
|
// is not one: a forged or garbled message is refused when its turn comes; here it is only named.
|
|
func declaredIn(body []byte) string {
|
|
var signed Signed
|
|
if err := json.Unmarshal(body, &signed); err != nil || len(signed.Declaration) == 0 {
|
|
return ""
|
|
}
|
|
sum := sha256.Sum256(signed.Declaration)
|
|
return hex.EncodeToString(sum[:])
|
|
}
|
|
|
|
// SayHealth says a health statement on the link open now (novox/hq ADR 0240, to-be 48 §4), and whether
|
|
// the bus took it. **Not through the worker**: a statement is a word about the machine as it is, not an
|
|
// account of an apply, and it must not wait behind one — a container crash-looping while a long apply
|
|
// runs is said while it runs. False with no link open, or one that cannot say it; the caller says it
|
|
// again on its next look, and the next report carries it whatever happens.
|
|
func (q *Queue) SayHealth(ctx context.Context, h Health) bool {
|
|
q.init()
|
|
q.mu.Lock()
|
|
bus := q.bus
|
|
q.mu.Unlock()
|
|
sayer, ok := bus.(HealthBus)
|
|
if bus == nil || !ok {
|
|
return false
|
|
}
|
|
body, err := json.Marshal(HealthSaid{Node: q.Membership.Node, Health: h})
|
|
if err != nil {
|
|
return false
|
|
}
|
|
saying, cancel := context.WithTimeout(ctx, q.timeout())
|
|
defer cancel()
|
|
if err := sayer.Health(saying, q.Membership.Node, body); err != nil {
|
|
q.Say("could not say this machine's health: " + err.Error())
|
|
return false
|
|
}
|
|
return true
|
|
}
|
|
|
|
func (q *Queue) timeout() time.Duration {
|
|
if q.Timeout <= 0 {
|
|
return 10 * time.Second
|
|
}
|
|
return q.Timeout
|
|
}
|
|
|
|
// AskTool asks a module's tool on this machine's own node tools, as a declared health check (novox/hq ADR
|
|
// 0240, to-be 48 §3), and reads its answer: healthy or not, with why. An error when there is no link open
|
|
// or nothing answered — said as the check's failure, never as healthy.
|
|
func (q *Queue) AskTool(ctx context.Context, module, tool string) (bool, string, error) {
|
|
q.init()
|
|
q.mu.Lock()
|
|
bus := q.bus
|
|
q.mu.Unlock()
|
|
asker, ok := bus.(ToolBus)
|
|
if bus == nil || !ok {
|
|
return false, "", errors.New("no link to the bus is open")
|
|
}
|
|
reply, err := asker.Ask(ctx, ToolSubject(module, tool, q.Membership.Node), []byte("{}"))
|
|
if err != nil {
|
|
return false, "", err
|
|
}
|
|
return ReadToolAnswer(reply)
|
|
}
|
|
|
|
// ReadToolAnswer reads a tool's reply envelope — `{result, error}`, as the node tools answer every call —
|
|
// as a health answer: an error is not healthy, and so is a result that does not say `healthy: true`.
|
|
func ReadToolAnswer(reply []byte) (bool, string, error) {
|
|
var env struct {
|
|
Result json.RawMessage `json:"result"`
|
|
Error string `json:"error"`
|
|
}
|
|
if err := json.Unmarshal(reply, &env); err != nil {
|
|
return false, "", fmt.Errorf("its answer is not one: %w", err)
|
|
}
|
|
if env.Error != "" {
|
|
return false, env.Error, nil
|
|
}
|
|
var a struct {
|
|
Healthy bool `json:"healthy"`
|
|
Why string `json:"why"`
|
|
}
|
|
if err := json.Unmarshal(env.Result, &a); err != nil {
|
|
return false, "it answered something other than healthy or not", nil
|
|
}
|
|
return a.Healthy, a.Why, nil
|
|
}
|