Merge pull request 'Say again what was last applied when the mesh asks (hq to-be 45 Phase 3, the report verb)' (#37) from feat/a-core-that-cannot-fail-silently-phase-3 into main
mesh/delivery held for a person: merged without a passing check: only a person decides that it goes on
mesh/delivery held for a person: merged without a passing check: only a person decides that it goes on
This commit was merged in pull request #37.
This commit is contained in:
@@ -161,7 +161,7 @@
|
||||
"type": "file",
|
||||
"path": "/var/lib/mesh-bus-conf/accounts.conf",
|
||||
"mode": "0600",
|
||||
"content": "// The first user list, carried by the installer because at genesis there is no mesh to\n// compose one. A bootstrap credential, rotated with the store's and replaced by the\n// controller's own composition from its first start onward.\naccounts {\n MESH {\n jetstream: enabled\n users = [\n { user: \"controller\", password: \"$2a$10$AHqJgOifIVbU41KmATiMhuXFs8xa7Wl2HuN4UVBCXdN2jIQzjqApy\", permissions: {\n publish: { allow: [\"$JS.ACK.CONTROL.controller.>\", \"$JS.ACK.EVENTS.controller.>\", \"$JS.API.>\", \"$KV.mesh-controller_calls.>\", \"$KV.mesh-controller_hand-acts.>\", \"$KV.mesh-controller_conditions.>\", \"$KV.mesh-controller_condition-history.>\", \"$KV.mesh-controller_lease.>\", \"$KV.SEAT_MESH_BUILD_MACHINE_cancelled.>\", \"$KV.SEAT_NODE_BUILD_AGENT_cancelled.>\", \"_INBOX.enrol.>\", \"mesh.assignment.>\", \"mesh.mod.*.tool.>\", \"mesh.node.>\", \"mesh.seat.mesh-build-machine.accept.>\", \"mesh.seat.node-build-agent.accept.>\", \"mesh.seat.mesh-build-machine.tool.>\", \"mesh.seat.node-build-agent.tool.>\", \"mesh.seat.mesh-controller.event.applied\", \"mesh.seat.mesh-controller.event.built-before\", \"mesh.seat.mesh-controller.event.refused\", \"mesh.seat.mesh-controller.event.condition-raised\", \"mesh.seat.mesh-controller.event.condition-changed\", \"mesh.seat.mesh-controller.event.condition-cleared\", \"mesh.seat.mesh-controller.event.doctor-heartbeat\", \"mesh.seat.mesh-controller.event.secret-replaced\", \"$SRV.INFO\", \"mesh.seat.node-intrusion-prevention.tool.banned.*\"] }\n subscribe: { allow: [\"$JS.API.>\", \"$JS.EVENT.ADVISORY.CONSUMER.MAX_DELIVERIES.>\", \"$JS.EVENT.ADVISORY.CONSUMER.DELETED.>\", \"_DELIVER.controller\", \"_DELIVER.controller.>\", \"_INBOX.controller.>\", \"mesh.control.>\", \"mesh.mod.*.event.provisioner.failing\", \"mesh.mod.*.event.provisioner.recovered\", \"mesh.mod.gitea.event.pull.merged\", \"mesh.mod.mesh-catalog.event.catching-up\", \"mesh.mod.mesh-catalog.event.upgraded\", \"mesh.seat.mesh-build-machine.event.built\", \"mesh.seat.node-build-agent.event.built\", \"mesh.seat.mesh-controller.tool.>\", \"$SRV.PING\", \"$SRV.INFO\", \"$SRV.PING.mesh-controller\", \"$SRV.PING.mesh-controller.>\", \"$SRV.INFO.mesh-controller\", \"$SRV.INFO.mesh-controller.>\", \"$SRV.STATS\", \"$SRV.STATS.mesh-controller\", \"$SRV.STATS.mesh-controller.>\"] }\n allow_responses: { max: 1, ttl: \"1m\" }\n } }\n ]\n }\n}\n"
|
||||
"content": "// The first user list, carried by the installer because at genesis there is no mesh to\n// compose one. A bootstrap credential, rotated with the store's and replaced by the\n// controller's own composition from its first start onward.\naccounts {\n MESH {\n jetstream: enabled\n users = [\n { user: \"controller\", password: \"$2a$10$AHqJgOifIVbU41KmATiMhuXFs8xa7Wl2HuN4UVBCXdN2jIQzjqApy\", permissions: {\n publish: { allow: [\"$JS.ACK.CONTROL.controller.>\", \"$JS.ACK.EVENTS.controller.>\", \"$JS.API.>\", \"$KV.mesh-controller_calls.>\", \"$KV.mesh-controller_hand-acts.>\", \"$KV.mesh-controller_conditions.>\", \"$KV.mesh-controller_condition-history.>\", \"$KV.mesh-controller_lease.>\", \"$KV.SEAT_MESH_BUILD_MACHINE_cancelled.>\", \"$KV.SEAT_NODE_BUILD_AGENT_cancelled.>\", \"_INBOX.enrol.>\", \"mesh.assignment.>\", \"mesh.mod.*.tool.>\", \"mesh.node.>\", \"mesh.seat.mesh-build-machine.accept.>\", \"mesh.seat.node-build-agent.accept.>\", \"mesh.seat.mesh-build-machine.tool.>\", \"mesh.seat.node-build-agent.tool.>\", \"mesh.seat.mesh-controller.event.applied\", \"mesh.seat.mesh-controller.event.built-before\", \"mesh.seat.mesh-controller.event.refused\", \"mesh.seat.mesh-controller.event.condition-raised\", \"mesh.seat.mesh-controller.event.condition-changed\", \"mesh.seat.mesh-controller.event.condition-cleared\", \"mesh.seat.mesh-controller.event.doctor-heartbeat\", \"mesh.seat.mesh-controller.event.secret-replaced\", \"mesh.seat.mesh-controller.event.healer-acted\", \"$SRV.INFO\", \"mesh.seat.node-intrusion-prevention.tool.banned.*\"] }\n subscribe: { allow: [\"$JS.API.>\", \"$JS.EVENT.ADVISORY.CONSUMER.MAX_DELIVERIES.>\", \"$JS.EVENT.ADVISORY.CONSUMER.DELETED.>\", \"_DELIVER.controller\", \"_DELIVER.controller.>\", \"_INBOX.controller.>\", \"mesh.control.>\", \"mesh.mod.*.event.provisioner.failing\", \"mesh.mod.*.event.provisioner.recovered\", \"mesh.mod.gitea.event.pull.merged\", \"mesh.mod.mesh-catalog.event.catching-up\", \"mesh.mod.mesh-catalog.event.upgraded\", \"mesh.seat.mesh-build-machine.event.built\", \"mesh.seat.node-build-agent.event.built\", \"mesh.seat.mesh-controller.tool.>\", \"$SRV.PING\", \"$SRV.INFO\", \"$SRV.PING.mesh-controller\", \"$SRV.PING.mesh-controller.>\", \"$SRV.INFO.mesh-controller\", \"$SRV.INFO.mesh-controller.>\", \"$SRV.STATS\", \"$SRV.STATS.mesh-controller\", \"$SRV.STATS.mesh-controller.>\"] }\n allow_responses: { max: 1, ttl: \"1m\" }\n } }\n ]\n }\n}\n"
|
||||
},
|
||||
{
|
||||
"id": "broker",
|
||||
|
||||
@@ -0,0 +1,130 @@
|
||||
package link
|
||||
|
||||
import (
|
||||
"context"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
// The `report` verb (novox/hq to-be 45 §6), what healer H1 asks a machine whose send went unreported:
|
||||
// the node says again what it last applied, as an ordinary report through the one queue.
|
||||
|
||||
// **Asked, the reconcile's account is said even when it is not news** — the mesh asked because it does
|
||||
// not know, and silence would be the very fault it is asking about.
|
||||
func TestAskedToReportSaysWhatWasLastAppliedEvenWhenNotNews(t *testing.T) {
|
||||
m, _ := aMember(t)
|
||||
l := newQuietLink()
|
||||
var lines []string
|
||||
q := aQueue(m, (&applies{}).apply, &keptInMemory{})
|
||||
q.Say = func(s string) { lines = append(lines, s) }
|
||||
q.Reconcile = func(context.Context) (Report, bool) {
|
||||
return Report{Declared: "d-kept", Order: Order{Epoch: 7, Sequence: 12}, Applied: []string{"a"}}, true
|
||||
}
|
||||
q.News = func(Report) bool { return false }
|
||||
q.attach(context.Background(), l)
|
||||
ctx, stop := context.WithCancel(context.Background())
|
||||
worker := runWorker(ctx, q)
|
||||
|
||||
q.ReconcileDue() // an ordinary reconcile: not news, not said
|
||||
time.Sleep(50 * time.Millisecond)
|
||||
if n := len(l.said()); n != 0 {
|
||||
t.Fatalf("an ordinary reconcile that is not news said %d report(s)", n)
|
||||
}
|
||||
q.ReportAsked()
|
||||
eventually(t, "asked, the node said nothing", func() bool { return len(l.said()) == 1 })
|
||||
stop()
|
||||
<-worker
|
||||
r := l.said()[0]
|
||||
if r.Declared != "d-kept" || r.Sequence != 12 || r.Epoch != 7 || r.ReportSequence == 0 || r.Node != m.Node {
|
||||
t.Fatalf("the account said is not the one kept, in order: %+v", r)
|
||||
}
|
||||
if !strings.Contains(strings.Join(lines, "\n"), "asked by the mesh") {
|
||||
t.Errorf("being asked is not said in the journal: %q", lines)
|
||||
}
|
||||
}
|
||||
|
||||
// **A delivery waiting answers the question better**: it is applied, once, and its report is the answer.
|
||||
func TestAskedWhileADeliveryWaitsIsThatDeliverysReport(t *testing.T) {
|
||||
m, key := aMember(t)
|
||||
l := newQuietLink()
|
||||
a := &applies{}
|
||||
q := aQueue(m, a.apply, &keptInMemory{})
|
||||
reconciled := 0
|
||||
q.Reconcile = func(context.Context) (Report, bool) { reconciled++; return Report{}, true }
|
||||
q.attach(context.Background(), l)
|
||||
q.Deliver(ordered(t, key, "new", 7, 13))
|
||||
q.ReportAsked()
|
||||
ctx, stop := context.WithCancel(context.Background())
|
||||
worker := runWorker(ctx, q)
|
||||
eventually(t, "the delivery was never reported", func() bool { return len(l.said()) == 1 })
|
||||
time.Sleep(50 * time.Millisecond)
|
||||
stop()
|
||||
<-worker
|
||||
if reconciled != 0 || len(l.said()) != 1 || l.said()[0].Sequence != 13 {
|
||||
t.Fatalf("reconciled %d, said %+v", reconciled, l.said())
|
||||
}
|
||||
}
|
||||
|
||||
// askingLink is a quiet link the mesh can ask.
|
||||
type askingLink struct {
|
||||
*quietLink
|
||||
asks chan struct{}
|
||||
}
|
||||
|
||||
func (l askingLink) AskedToReport() <-chan struct{} { return l.asks }
|
||||
|
||||
// **The link hands the question to the queue**: asked on the link, the account is said on it.
|
||||
func TestALinkThatIsAskedHandsItToTheQueue(t *testing.T) {
|
||||
m, _ := aMember(t)
|
||||
l := askingLink{quietLink: newQuietLink(), asks: make(chan struct{}, 1)}
|
||||
q := aQueue(m, (&applies{}).apply, &keptInMemory{})
|
||||
q.Reconcile = func(context.Context) (Report, bool) { return Report{Declared: "d-kept"}, true }
|
||||
ctx, stop := context.WithCancel(context.Background())
|
||||
defer stop()
|
||||
worker := runWorker(ctx, q)
|
||||
served := make(chan error, 1)
|
||||
go func() { served <- serve(ctx, l, m, q, nil, time.Second) }()
|
||||
l.asks <- struct{}{}
|
||||
eventually(t, "asked on the link, nothing was said", func() bool {
|
||||
for _, r := range l.said() {
|
||||
if r.Declared == "d-kept" {
|
||||
return true
|
||||
}
|
||||
}
|
||||
return false
|
||||
})
|
||||
stop()
|
||||
<-served
|
||||
<-worker
|
||||
}
|
||||
|
||||
// **Against a real server**: the question on the node's own subject reaches it, and another's does not.
|
||||
func TestNatsTheMeshAsksANodeToReport(t *testing.T) {
|
||||
conn, js := aBus(t)
|
||||
l := &natsLink{conn: conn, js: js, node: "asked"}
|
||||
l.hearAsks()
|
||||
defer l.Close()
|
||||
if err := conn.Publish(AskReportSubject("another"), []byte(`{}`)); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := conn.Publish(AskReportSubject("asked"), []byte(`{"by":"the controller's healer H1"}`)); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := conn.Flush(); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
select {
|
||||
case <-l.AskedToReport():
|
||||
case <-time.After(5 * time.Second):
|
||||
t.Fatal("the question on the node's own subject never reached it")
|
||||
}
|
||||
select {
|
||||
case <-l.AskedToReport():
|
||||
t.Fatal("asked twice: another node's question reached this one")
|
||||
case <-time.After(200 * time.Millisecond):
|
||||
}
|
||||
if AskReportSubject("asked") != "mesh.node.asked.ask.report" {
|
||||
t.Fatalf("the subject %q is not the one the controller asks on and grants", AskReportSubject("asked"))
|
||||
}
|
||||
}
|
||||
@@ -38,6 +38,13 @@ type Link interface {
|
||||
Close()
|
||||
}
|
||||
|
||||
// Asked is a link that hears the mesh asking this node to say again what it last applied (novox/hq
|
||||
// to-be 45 §6, the `report` verb). Optional: a link without it is never asked, and the mesh then sends
|
||||
// the current declaration again instead.
|
||||
type Asked interface {
|
||||
AskedToReport() <-chan struct{}
|
||||
}
|
||||
|
||||
// Declaration is one thing the mesh told this node to be.
|
||||
//
|
||||
// **Handled, once — after the report is published.** A node that dies between applying and
|
||||
|
||||
@@ -38,6 +38,11 @@ const EnrolSubject = "mesh.control.enrol"
|
||||
// for.
|
||||
func DeclareSubject(node string) string { return "mesh.node." + node + ".declare" }
|
||||
|
||||
// AskReportSubject is where the mesh asks this node to say again what it last applied (novox/hq to-be 45
|
||||
// §6): its own, on core NATS and off any stream. The answer is an ordinary report on its report subject,
|
||||
// the one thing it already says — nobody's inbox is answered.
|
||||
func AskReportSubject(node string) string { return "mesh.node." + node + ".ask.report" }
|
||||
|
||||
// natsURL is a bus address as the client wants it. A membership records host and port, because that
|
||||
// is what genesis sealed into it and what the other transport takes; the scheme is this transport's
|
||||
// own business.
|
||||
@@ -53,6 +58,8 @@ type natsLink struct {
|
||||
conn *nats.Conn
|
||||
js nats.JetStreamContext
|
||||
sub *nats.Subscription
|
||||
asking *nats.Subscription
|
||||
asked chan struct{}
|
||||
node string
|
||||
arrived chan Declaration
|
||||
lost chan error
|
||||
@@ -106,6 +113,7 @@ func dialNats(ctx context.Context, m Membership, timeout time.Duration) (Link, e
|
||||
conn: conn, js: js, node: m.Node,
|
||||
arrived: make(chan Declaration, drainDepth),
|
||||
lost: make(chan error, 1),
|
||||
asked: make(chan struct{}, 1),
|
||||
}
|
||||
|
||||
// Bound to the consumer the controller made for this node, named after the node because that is
|
||||
@@ -124,6 +132,8 @@ func dialNats(ctx context.Context, m Membership, timeout time.Duration) (Link, e
|
||||
}
|
||||
l.sub = sub
|
||||
|
||||
l.hearAsks()
|
||||
|
||||
conn.SetDisconnectErrHandler(func(_ *nats.Conn, err error) {
|
||||
select {
|
||||
case l.lost <- fmt.Errorf("the link dropped: %w", err):
|
||||
@@ -163,13 +173,36 @@ func dialNats(ctx context.Context, m Membership, timeout time.Duration) (Link, e
|
||||
return l, nil
|
||||
}
|
||||
|
||||
// hearAsks listens for the mesh asking what this node last applied (the `report` verb). One question
|
||||
// waiting is enough: a second asked before the first is taken is the same question. A bus whose user
|
||||
// list is older than the grant refuses the subscription asynchronously — said by the client, and the
|
||||
// mesh then sends again instead of asking; nothing here depends on being asked.
|
||||
func (l *natsLink) hearAsks() {
|
||||
if l.asked == nil {
|
||||
l.asked = make(chan struct{}, 1)
|
||||
}
|
||||
asking, err := l.conn.Subscribe(AskReportSubject(l.node), func(*nats.Msg) {
|
||||
select {
|
||||
case l.asked <- struct{}{}:
|
||||
default:
|
||||
}
|
||||
})
|
||||
if err == nil {
|
||||
l.asking = asking
|
||||
}
|
||||
}
|
||||
|
||||
func (l *natsLink) Declarations() <-chan Declaration { return l.arrived }
|
||||
func (l *natsLink) AskedToReport() <-chan struct{} { return l.asked }
|
||||
func (l *natsLink) Lost() <-chan error { return l.lost }
|
||||
|
||||
func (l *natsLink) Close() {
|
||||
if l.sub != nil {
|
||||
_ = l.sub.Unsubscribe()
|
||||
}
|
||||
if l.asking != nil {
|
||||
_ = l.asking.Unsubscribe()
|
||||
}
|
||||
if l.conn != nil {
|
||||
l.conn.Close()
|
||||
}
|
||||
|
||||
+36
-7
@@ -64,8 +64,11 @@ type Queue struct {
|
||||
delivering bool
|
||||
waiting []Declaration // delivered and not yet taken, in the order they arrived
|
||||
due bool // a reconcile asked for
|
||||
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
|
||||
// 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.
|
||||
@@ -116,6 +119,19 @@ func (q *Queue) ReconcileDue() {
|
||||
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{}{}:
|
||||
@@ -147,8 +163,8 @@ func (q *Queue) step(ctx context.Context) bool {
|
||||
q.mu.Unlock()
|
||||
return false
|
||||
}
|
||||
batch, due := q.waiting, q.due
|
||||
q.waiting, q.due = nil, 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() {
|
||||
@@ -159,7 +175,7 @@ func (q *Queue) step(ctx context.Context) bool {
|
||||
}()
|
||||
|
||||
if len(batch) == 0 {
|
||||
q.reconcile(ctx)
|
||||
q.reconcile(ctx, asked)
|
||||
return true
|
||||
}
|
||||
|
||||
@@ -210,14 +226,27 @@ func (q *Queue) step(ctx context.Context) bool {
|
||||
return true
|
||||
}
|
||||
|
||||
func (q *Queue) reconcile(ctx context.Context) {
|
||||
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 || q.News == nil || !q.News(report) {
|
||||
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)
|
||||
}
|
||||
|
||||
|
||||
@@ -234,6 +234,11 @@ func serve(ctx context.Context, link Link, m Membership, queue *Queue, say Annou
|
||||
defer queue.detach(link)
|
||||
|
||||
declarations := link.Declarations()
|
||||
// And the mesh asking what this node last applied (to-be 45 §6), on a link that hears the question.
|
||||
var asks <-chan struct{}
|
||||
if a, ok := link.(Asked); ok {
|
||||
asks = a.AskedToReport()
|
||||
}
|
||||
|
||||
for {
|
||||
// Asked to stop — by standing aside, say — nothing further is taken in.
|
||||
@@ -247,6 +252,8 @@ func serve(ctx context.Context, link Link, m Membership, queue *Queue, say Annou
|
||||
publishAlive(ctx, link, m, say, timeout)
|
||||
case reason := <-link.Lost():
|
||||
return reason
|
||||
case <-asks:
|
||||
queue.ReportAsked()
|
||||
case declaration, ok := <-declarations:
|
||||
if !ok {
|
||||
// The link's own reason, when it has managed to say one: "stopped delivering" on
|
||||
|
||||
Reference in New Issue
Block a user