The controller's healer H1 answers a send that went unreported by asking the machine first: a report lost on its way (issue 264) needs no second send. The node-engine now hears mesh.node.<self>.ask.report on core NATS and enqueues a reconcile whose account is said whether or not it is news; a delivery waiting meanwhile is applied and reported instead. The answer is an ordinary report on its own subject, so the node publishes nothing new and answers nobody's inbox. The genesis lock grants the controller the healer-acted seat event, which it now composes; the genesis test in mesh-controller holds the two equal.
131 lines
4.4 KiB
Go
131 lines
4.4 KiB
Go
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"))
|
|
}
|
|
}
|