The host's inbound behind a seam, with both transports
The outbound half went behind `Bus` and a node's two statements stopped naming a transport. This is the other half, and where the transport reached furthest: the run loop selected on a channel of the client library's own delivery type, so every part of holding a node in its mesh knew which bus it was on. `Link` is dialling, hearing and saying in one interface, because dialling is where the transport is chosen and choosing it twice is how one half of a node ends up on a different bus from the other. `Declaration` has one way of being done rather than two: a declaration set aside for a newer one is settled exactly as an applied one is, on both buses, and the difference is a fact the report carries. Four things this settled. **The host declares nothing on the new bus.** On the bus the mesh has it declares its own queue, because a queue that is not there means a node that hears nothing. Here it binds to a consumer the mesh made when the node enrolled, and a missing one is said as the mesh's to answer rather than quietly created with whatever this client happens to default to. **The pin is easier here than in the tool runtime, not harder.** The Go client takes a *tls.Config, so the same PinnedConfig with the same VerifyPeerCertificate does the work — the subject-alternative-name constraint recorded against the runtime's client is that client's, because it takes PEM strings with no verify hook. A host checks the fingerprint and nothing else. **Binding needs the subject as well as the consumer.** An empty subject is refused rather than taken to mean "whatever that consumer delivers", which the server said plainly and only when asked. **Reconnection stays the caller's.** Hold already decides when to try again and how long to wait; a client reconnecting underneath it would make that reasoning a duplicate of the library's. The drain keeps its live half and loses its catch-up half, as it said it would: verified that three declarations pushed to an absent node leave one on the stream, and it is the newest. One test-harness lesson worth the comment it got: delete-then-add is not a reset. A test that did that inherited the previous test's messages, and the symptom was a declaration counted as delivered twice — which reads as a redelivery bug in the code under test rather than as a dirty stream.
This commit is contained in:
@@ -0,0 +1,84 @@
|
||||
package link
|
||||
|
||||
import (
|
||||
"context"
|
||||
"time"
|
||||
)
|
||||
|
||||
// What a host hears, as the host's own words for it.
|
||||
//
|
||||
// The outbound half went behind `Bus` (bus.go) and a node's two statements stopped naming a
|
||||
// transport. This is the other half — dialling, and the declarations that arrive — and it is where
|
||||
// the transport reached furthest: the run loop selected on a channel of the client library's own
|
||||
// delivery type, so every part of holding a node in its mesh knew which bus it was on.
|
||||
//
|
||||
// **The host still imports nothing of the mesh's own** (novox/hq ADR 0005). This is its own
|
||||
// interface over its own libraries, and it agrees with the controller only because a conformance
|
||||
// fixture holds both to one envelope.
|
||||
|
||||
// Link is this node's live connection to its mesh: what it hears, and what it says.
|
||||
//
|
||||
// One interface rather than two, because **dialling is where the transport is chosen** and choosing
|
||||
// it twice is how one half of a node ends up on a different bus from the other.
|
||||
type Link interface {
|
||||
// Bus is what this node says: what it applied, and that it is here.
|
||||
Bus
|
||||
|
||||
// Declarations is what the mesh tells this node to be.
|
||||
Declarations() <-chan Declaration
|
||||
|
||||
// Lost says the link ended, and why.
|
||||
//
|
||||
// **Read rather than discovered.** A node that finds out by noticing silence is a node that
|
||||
// believed it was in the mesh for as long as the silence lasted, which is the one state ADR
|
||||
// 0004 says must never look like being connected.
|
||||
Lost() <-chan error
|
||||
|
||||
// Close lets go of whatever was dialled.
|
||||
Close()
|
||||
}
|
||||
|
||||
// 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
|
||||
// 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.
|
||||
//
|
||||
// There is one way of being done rather than two. A declaration set aside because a newer arrived
|
||||
// with it is settled exactly as an applied one is, on both buses, and the difference between them
|
||||
// is a fact the *report* carries — a second method here would be a distinction the transport does
|
||||
// not make.
|
||||
type Declaration interface {
|
||||
// Body is the signed declaration as it arrived, bytes unchanged: a node verifies what it
|
||||
// received rather than what it re-encoded.
|
||||
Body() []byte
|
||||
|
||||
// Handled settles it. Called after the report for it has been published, either way.
|
||||
Handled() error
|
||||
}
|
||||
|
||||
// Open opens this node's link to its mesh.
|
||||
//
|
||||
// Named Open rather than Dial because Dial is this package's raw TLS dial, which the enrolment path
|
||||
// uses to see a certificate before it trusts anything.
|
||||
//
|
||||
// **Both transports ship and this is the one place that chooses** (novox/hq ADR 0116: nothing moves
|
||||
// a node's bus before step 5). Until then every membership names the bus the mesh runs on today,
|
||||
// and the rollout is this switch and the credential behind it — not a change anywhere in the loop
|
||||
// that reads from what comes back.
|
||||
func Open(ctx context.Context, m Membership, timeout time.Duration) (Link, error) {
|
||||
switch m.Transport {
|
||||
case OnNATS:
|
||||
return dialNats(ctx, m, timeout)
|
||||
default:
|
||||
return dialCurrent(ctx, m, timeout)
|
||||
}
|
||||
}
|
||||
|
||||
// The buses a node can be on. Empty is the one the mesh runs on today, which is every node until
|
||||
// the rollout — so a membership recorded before any of this existed reads as correct rather than as
|
||||
// unset.
|
||||
const (
|
||||
OnCurrent = ""
|
||||
OnNATS = "nats"
|
||||
)
|
||||
@@ -0,0 +1,164 @@
|
||||
package link
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"net/url"
|
||||
"time"
|
||||
|
||||
amqp "github.com/rabbitmq/amqp091-go"
|
||||
)
|
||||
|
||||
// The host's link on the bus the mesh runs on today.
|
||||
//
|
||||
// Moved out of the run loop rather than changed: the dial, the queue, the prefetch window and the
|
||||
// return handler are what they were, because the mesh is running on this and a bus nothing speaks
|
||||
// yet is no reason to alter the one every node is on.
|
||||
|
||||
// currentLink is this node's connection as a channel.
|
||||
type currentLink struct {
|
||||
conn *amqp.Connection
|
||||
channel *amqp.Channel
|
||||
arrived chan Declaration
|
||||
lost chan error
|
||||
}
|
||||
|
||||
func dialCurrent(ctx context.Context, m Membership, timeout time.Duration) (Link, error) {
|
||||
config, err := PinnedConfig(m.Fingerprint)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
dsn := fmt.Sprintf("amqps://%s:%s@%s/",
|
||||
url.QueryEscape(m.Node), url.QueryEscape(m.Password), m.Broker)
|
||||
conn, err := amqp.DialConfig(dsn, amqp.Config{
|
||||
TLSClientConfig: config,
|
||||
Dial: amqp.DefaultDial(timeout),
|
||||
// Kept short so a node that has silently lost its route notices, rather than holding a
|
||||
// connection the broker forgot about and believing it is still in the mesh.
|
||||
Heartbeat: 10 * time.Second,
|
||||
})
|
||||
if err != nil {
|
||||
if errors.Is(err, ErrWrongCertificate) {
|
||||
return nil, err
|
||||
}
|
||||
return nil, fmt.Errorf("cannot reach the broker at %s: %w", m.Broker, err)
|
||||
}
|
||||
|
||||
channel, err := conn.Channel()
|
||||
if err != nil {
|
||||
conn.Close()
|
||||
return nil, err
|
||||
}
|
||||
|
||||
queue := QueueFor(m.Node)
|
||||
if _, err := channel.QueueDeclare(queue, true, false, false, false, nil); err != nil {
|
||||
conn.Close()
|
||||
return nil, fmt.Errorf("cannot declare this node's queue %s: %w", queue, err)
|
||||
}
|
||||
|
||||
// Applying is one at a time — two at once would race on the same filesystem — but SEEING is
|
||||
// not: with a prefetch of one the host could never know that a newer declaration was already
|
||||
// waiting, and so applied every one of a backlog in turn, at the better part of a minute each,
|
||||
// becoming things nobody wanted any more (novox/hq issue 031). A window of unacknowledged
|
||||
// deliveries lets it drain to the newest; each declaration still survives a restart on the
|
||||
// broker until it is acknowledged, which happens only after it is applied or set aside.
|
||||
if err := channel.Qos(drainDepth, 0, false); err != nil {
|
||||
conn.Close()
|
||||
return nil, err
|
||||
}
|
||||
|
||||
deliveries, err := channel.ConsumeWithContext(ctx, queue, "", false, false, false, false, nil)
|
||||
if err != nil {
|
||||
conn.Close()
|
||||
return nil, err
|
||||
}
|
||||
|
||||
l := ¤tLink{
|
||||
conn: conn, channel: channel,
|
||||
arrived: make(chan Declaration, drainDepth),
|
||||
lost: make(chan error, 1),
|
||||
}
|
||||
|
||||
// Published mandatory, so the broker hands back anything it cannot route rather than dropping
|
||||
// it. Without this a report goes to an exchange with no matching binding, the publisher is told
|
||||
// nothing, and the mesh believes this node never answered while the node believes it did —
|
||||
// which is what happened when `report` was left unbound on the other side.
|
||||
returned := channel.NotifyReturn(make(chan amqp.Return, 4))
|
||||
go func() {
|
||||
for r := range returned {
|
||||
select {
|
||||
case l.lost <- fmt.Errorf("the broker could not route this node's %s: %s (%d %s)",
|
||||
r.RoutingKey, r.Exchange, r.ReplyCode, r.ReplyText):
|
||||
default:
|
||||
}
|
||||
}
|
||||
}()
|
||||
|
||||
closed := conn.NotifyClose(make(chan *amqp.Error, 1))
|
||||
go func() {
|
||||
select {
|
||||
case reason := <-closed:
|
||||
select {
|
||||
case l.lost <- fmt.Errorf("the link closed: %v", reason):
|
||||
default:
|
||||
}
|
||||
case <-ctx.Done():
|
||||
}
|
||||
}()
|
||||
|
||||
// One goroutine turning the library's deliveries into the mesh's words, so the run loop selects
|
||||
// on one kind of thing whichever bus it is on.
|
||||
go func() {
|
||||
defer close(l.arrived)
|
||||
for {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return
|
||||
case delivery, ok := <-deliveries:
|
||||
if !ok {
|
||||
select {
|
||||
case l.lost <- errors.New("the broker stopped delivering"):
|
||||
default:
|
||||
}
|
||||
return
|
||||
}
|
||||
select {
|
||||
case l.arrived <- currentDeclaration{delivery}:
|
||||
case <-ctx.Done():
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
}()
|
||||
|
||||
return l, nil
|
||||
}
|
||||
|
||||
func (l *currentLink) Declarations() <-chan Declaration { return l.arrived }
|
||||
func (l *currentLink) Lost() <-chan error { return l.lost }
|
||||
|
||||
func (l *currentLink) Close() {
|
||||
if l.channel != nil {
|
||||
_ = l.channel.Close()
|
||||
}
|
||||
if l.conn != nil {
|
||||
_ = l.conn.Close()
|
||||
}
|
||||
}
|
||||
|
||||
// Report and Alive are the outbound half, over the channel this link holds.
|
||||
func (l *currentLink) Report(ctx context.Context, node string, body []byte) error {
|
||||
return OverCurrent{Channel: l.channel}.Report(ctx, node, body)
|
||||
}
|
||||
|
||||
func (l *currentLink) Alive(ctx context.Context, node string, body []byte) error {
|
||||
return OverCurrent{Channel: l.channel}.Alive(ctx, node, body)
|
||||
}
|
||||
|
||||
// currentDeclaration is one delivery from the bus the mesh has.
|
||||
type currentDeclaration struct{ delivery amqp.Delivery }
|
||||
|
||||
func (d currentDeclaration) Body() []byte { return d.delivery.Body }
|
||||
func (d currentDeclaration) Handled() error { return d.delivery.Ack(false) }
|
||||
@@ -0,0 +1,182 @@
|
||||
package link
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/nats-io/nats.go"
|
||||
)
|
||||
|
||||
// The host's link on the bus being built.
|
||||
//
|
||||
// Two things a host does here that it cannot do on the other bus, and one it must not try.
|
||||
//
|
||||
// **It declares nothing.** On the bus the mesh has, a host declares its own queue on connecting,
|
||||
// because a queue that is not there means a node that hears nothing. Here the object it reads
|
||||
// through is a durable consumer, and a host's account reaches no part of the JetStream API — by
|
||||
// design, because the controller is the only writer of consumer definitions (design 25 §3). So the
|
||||
// host **binds** to a consumer the controller made when this node enrolled, and a missing one is
|
||||
// said as what it is rather than quietly created with whatever configuration this client happens to
|
||||
// default to.
|
||||
//
|
||||
// **It gets order for free, and keeps the drain anyway.** The declaration subject is last-per-subject
|
||||
// (design 29 §4), so a node that was away receives exactly the current declaration rather than a
|
||||
// queue of superseded ones, and the stream's sequence orders them definitively — the wire-level
|
||||
// answer to novox/hq issue 107. What the drain in run.go still answers is the live case: three
|
||||
// pushes to a *connected* node are three deliveries whatever the stream later retains.
|
||||
|
||||
// DeclareSubject is where this node's declaration lands. Its own, and no other node's: a host's
|
||||
// account subscribes exactly this and the subject is the authority on which node a declaration is
|
||||
// for.
|
||||
func DeclareSubject(node string) string { return "mesh.node." + node + ".declare" }
|
||||
|
||||
// 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.
|
||||
func natsURL(address string) string {
|
||||
if strings.Contains(address, "://") {
|
||||
return address
|
||||
}
|
||||
return "nats://" + address
|
||||
}
|
||||
|
||||
// natsLink is this node's connection as a JetStream subscription.
|
||||
type natsLink struct {
|
||||
conn *nats.Conn
|
||||
js nats.JetStreamContext
|
||||
sub *nats.Subscription
|
||||
node string
|
||||
arrived chan Declaration
|
||||
lost chan error
|
||||
}
|
||||
|
||||
func dialNats(ctx context.Context, m Membership, timeout time.Duration) (Link, error) {
|
||||
// Pinned exactly as the other transport is, and for once the Go client makes that easy: it
|
||||
// takes a *tls.Config, so the same PinnedConfig with the same VerifyPeerCertificate does the
|
||||
// work. **The constraint recorded against the tool runtime does not apply here** — that client
|
||||
// takes PEM strings with no verify hook, which is why the bus's certificate must carry a name
|
||||
// matching the address *modules* dial it by. A host checks the fingerprint and nothing else.
|
||||
config, err := PinnedConfig(m.Fingerprint)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
opts := []nats.Option{
|
||||
nats.Secure(config),
|
||||
nats.UserInfo(m.Node, m.Password),
|
||||
nats.Name("mesh-host/" + m.Node),
|
||||
nats.Timeout(timeout),
|
||||
// A node that has silently lost its route notices, rather than holding a connection the
|
||||
// server forgot about and believing it is still in the mesh.
|
||||
nats.PingInterval(10 * time.Second),
|
||||
nats.MaxPingsOutstanding(2),
|
||||
// Reconnection is the caller's: Hold already decides when to try again and how long to
|
||||
// wait, and a client quietly reconnecting underneath it would make that reasoning a
|
||||
// duplicate of the library's.
|
||||
nats.NoReconnect(),
|
||||
}
|
||||
|
||||
conn, err := nats.Connect(natsURL(m.Broker), opts...)
|
||||
if err != nil {
|
||||
if errors.Is(err, ErrWrongCertificate) {
|
||||
return nil, err
|
||||
}
|
||||
return nil, fmt.Errorf("cannot reach the bus at %s: %w", m.Broker, err)
|
||||
}
|
||||
js, err := conn.JetStream()
|
||||
if err != nil {
|
||||
conn.Close()
|
||||
return nil, fmt.Errorf("the bus at %s has no JetStream: %w", m.Broker, err)
|
||||
}
|
||||
|
||||
l := &natsLink{
|
||||
conn: conn, js: js, node: m.Node,
|
||||
arrived: make(chan Declaration, drainDepth),
|
||||
lost: make(chan error, 1),
|
||||
}
|
||||
|
||||
// Bound to the consumer the controller made for this node, named after the node because that is
|
||||
// what the node's own ack grant allows (`$JS.ACK.NODES.<node>.>`).
|
||||
feed := make(chan *nats.Msg, drainDepth)
|
||||
// The subject as well as the binding: the client checks what is asked for against the
|
||||
// consumer's own filter, and an empty subject is refused rather than taken to mean "whatever
|
||||
// that consumer delivers".
|
||||
sub, err := js.ChanSubscribe(DeclareSubject(m.Node), feed, nats.Bind("NODES", m.Node))
|
||||
if err != nil {
|
||||
conn.Close()
|
||||
return nil, fmt.Errorf(
|
||||
"this node cannot read its declarations: %w. The mesh creates that when a node enrols, "+
|
||||
"and a host may not create one itself — so this is the mesh's to answer, not this "+
|
||||
"machine's", err)
|
||||
}
|
||||
l.sub = sub
|
||||
|
||||
conn.SetDisconnectErrHandler(func(_ *nats.Conn, err error) {
|
||||
select {
|
||||
case l.lost <- fmt.Errorf("the link dropped: %w", err):
|
||||
default:
|
||||
}
|
||||
})
|
||||
conn.SetClosedHandler(func(*nats.Conn) {
|
||||
select {
|
||||
case l.lost <- errors.New("the link closed"):
|
||||
default:
|
||||
}
|
||||
})
|
||||
|
||||
go func() {
|
||||
defer close(l.arrived)
|
||||
for {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return
|
||||
case msg, ok := <-feed:
|
||||
if !ok {
|
||||
select {
|
||||
case l.lost <- errors.New("the bus stopped delivering"):
|
||||
default:
|
||||
}
|
||||
return
|
||||
}
|
||||
select {
|
||||
case l.arrived <- natsDeclaration{msg}:
|
||||
case <-ctx.Done():
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
}()
|
||||
|
||||
return l, nil
|
||||
}
|
||||
|
||||
func (l *natsLink) Declarations() <-chan Declaration { return l.arrived }
|
||||
func (l *natsLink) Lost() <-chan error { return l.lost }
|
||||
|
||||
func (l *natsLink) Close() {
|
||||
if l.sub != nil {
|
||||
_ = l.sub.Unsubscribe()
|
||||
}
|
||||
if l.conn != nil {
|
||||
l.conn.Close()
|
||||
}
|
||||
}
|
||||
|
||||
func (l *natsLink) Report(ctx context.Context, node string, body []byte) error {
|
||||
return OverNATS{Conn: l.conn, JS: l.js}.Report(ctx, node, body)
|
||||
}
|
||||
|
||||
func (l *natsLink) Alive(ctx context.Context, node string, body []byte) error {
|
||||
return OverNATS{Conn: l.conn, JS: l.js}.Alive(ctx, node, body)
|
||||
}
|
||||
|
||||
// natsDeclaration is one declaration off the NODES stream.
|
||||
type natsDeclaration struct{ msg *nats.Msg }
|
||||
|
||||
func (d natsDeclaration) Body() []byte { return d.msg.Data }
|
||||
|
||||
// Handled acknowledges it. The ack goes to this node's own ack subject, which is the one thing
|
||||
// besides its reports a node's account may publish.
|
||||
func (d natsDeclaration) Handled() error { return d.msg.Ack() }
|
||||
@@ -0,0 +1,228 @@
|
||||
package link
|
||||
|
||||
import (
|
||||
"context"
|
||||
"crypto/ed25519"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"os"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/nats-io/nats.go"
|
||||
)
|
||||
|
||||
// The host's link against a real server, because every claim here is about one.
|
||||
//
|
||||
// Whether binding to a consumer the host did not create works, whether a declaration on the node's
|
||||
// own subject arrives, whether acknowledging it removes it from the consumer's pending — none of
|
||||
// that can be reasoned out, and the first two are the ones that would leave a node silently hearing
|
||||
// nothing:
|
||||
//
|
||||
// docker run -d --rm --name t -p 14223:4222 nats:2.10-alpine -js
|
||||
// MESH_TEST_NATS=nats://127.0.0.1:14223 go test ./internal/link/ -run TestNats
|
||||
|
||||
func aBus(t *testing.T) (*nats.Conn, nats.JetStreamContext) {
|
||||
t.Helper()
|
||||
url := os.Getenv("MESH_TEST_NATS")
|
||||
if url == "" {
|
||||
t.Skip("MESH_TEST_NATS unset")
|
||||
}
|
||||
conn, err := nats.Connect(url)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
t.Cleanup(conn.Close)
|
||||
js, err := conn.JetStream()
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
// **Ensured and purged, not deleted and recreated.** Delete-then-add looked like a reset and is
|
||||
// not one: a test that did that inherited the previous test's messages, and the symptom was a
|
||||
// declaration counted as delivered twice — which reads as a redelivery bug in the code under
|
||||
// test rather than as a dirty stream. Purge is defined to empty a stream; recreating one is a
|
||||
// race with the server's own teardown.
|
||||
for _, want := range []*nats.StreamConfig{
|
||||
{Name: "NODES", Subjects: []string{"mesh.node.*.declare"}, MaxMsgsPerSubject: 1},
|
||||
{Name: "CONTROL", Subjects: []string{"mesh.control.*.report", "mesh.control.enrol"},
|
||||
Retention: nats.WorkQueuePolicy},
|
||||
} {
|
||||
if _, err := js.StreamInfo(want.Name); err != nil {
|
||||
if _, err := js.AddStream(want); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
if err := js.PurgeStream(want.Name); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
return conn, js
|
||||
}
|
||||
|
||||
// theMeshMakes is the consumer the controller creates when a node enrols. Made here by the test
|
||||
// because the host may not: its account reaches no part of the JetStream API, which is the whole
|
||||
// reason this binds rather than subscribes.
|
||||
//
|
||||
// Removed afterwards, and each test names its own node: two tests sharing a consumer name share its
|
||||
// delivery count and its pending list, and the first thing that goes wrong reads as a fault in the
|
||||
// host rather than in the test beside it.
|
||||
func theMeshMakes(t *testing.T, js nats.JetStreamContext, node string) {
|
||||
t.Helper()
|
||||
t.Cleanup(func() { _ = js.DeleteConsumer("NODES", node) })
|
||||
if _, err := js.AddConsumer("NODES", &nats.ConsumerConfig{
|
||||
Durable: node,
|
||||
FilterSubject: DeclareSubject(node),
|
||||
AckPolicy: nats.AckExplicitPolicy,
|
||||
AckWait: 300 * time.Second,
|
||||
DeliverSubject: "_DELIVER." + node,
|
||||
}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
|
||||
func signedBy(t *testing.T, key ed25519.PrivateKey, declaration []byte) []byte {
|
||||
t.Helper()
|
||||
body, err := json.Marshal(Signed{
|
||||
Declaration: declaration, Signature: ed25519.Sign(key, declaration),
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
return body
|
||||
}
|
||||
|
||||
// A declaration on this node's own subject reaches the host, is applied, and acknowledging it
|
||||
// empties the consumer — which is what tells the mesh the node has it.
|
||||
func TestNatsADeclarationReachesTheHostAndIsSettled(t *testing.T) {
|
||||
conn, js := aBus(t)
|
||||
const node = "settling"
|
||||
theMeshMakes(t, js, node)
|
||||
|
||||
public, private, _ := ed25519.GenerateKey(nil)
|
||||
m := Membership{Node: node, Signer: public}
|
||||
|
||||
// Dialled directly rather than through Open: the test server has no TLS, and what is being
|
||||
// checked is the subscription and the settling, not the pin — which PinnedConfig owns and its
|
||||
// own tests cover.
|
||||
l := &natsLink{conn: conn, js: js, node: node,
|
||||
arrived: make(chan Declaration, drainDepth), lost: make(chan error, 1)}
|
||||
feed := make(chan *nats.Msg, drainDepth)
|
||||
sub, err := js.ChanSubscribe(DeclareSubject(node), feed, nats.Bind("NODES", node))
|
||||
if err != nil {
|
||||
t.Fatalf("the host could not bind to the consumer the mesh made for it: %v", err)
|
||||
}
|
||||
defer func() { _ = sub.Unsubscribe() }()
|
||||
go func() {
|
||||
for msg := range feed {
|
||||
l.arrived <- natsDeclaration{msg}
|
||||
}
|
||||
}()
|
||||
|
||||
if _, err := js.Publish(DeclareSubject(node),
|
||||
signedBy(t, private, []byte(`{"declared":"d1"}`))); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
select {
|
||||
case d := <-l.Declarations():
|
||||
report := handleBody(context.Background(), m, d.Body(),
|
||||
func(context.Context, []byte, []byte) Report {
|
||||
return Report{Applied: []string{"store"}}
|
||||
})
|
||||
if report.Refused != "" {
|
||||
t.Fatalf("a declaration the mesh signed was refused: %s", report.Refused)
|
||||
}
|
||||
if err := d.Handled(); err != nil {
|
||||
t.Fatalf("the node could not acknowledge its own declaration: %v", err)
|
||||
}
|
||||
case <-time.After(8 * time.Second):
|
||||
t.Fatal("no declaration reached the host")
|
||||
}
|
||||
|
||||
// **Nothing pending is the property**; a delivery count is not. Delivery is at-least-once by
|
||||
// design, so pinning "delivered exactly once" would be asserting something the mesh does not
|
||||
// rely on. What matters is that the acknowledgement landed, so the mesh can tell the node has
|
||||
// it — and that no redelivery was needed to get there, which is what would say the node was
|
||||
// too slow to answer for its own ack wait.
|
||||
deadline := time.Now().Add(5 * time.Second)
|
||||
var last string
|
||||
for time.Now().Before(deadline) {
|
||||
info, err := js.ConsumerInfo("NODES", node)
|
||||
switch {
|
||||
case err != nil:
|
||||
last = err.Error()
|
||||
case info.NumAckPending == 0 && info.NumRedelivered == 0:
|
||||
return
|
||||
default:
|
||||
last = fmt.Sprintf("pending %d, redelivered %d", info.NumAckPending, info.NumRedelivered)
|
||||
}
|
||||
time.Sleep(20 * time.Millisecond)
|
||||
}
|
||||
t.Fatalf("the declaration was not settled, so the mesh cannot tell the node has it: %s", last)
|
||||
}
|
||||
|
||||
// **A node that was away gets exactly the current declaration and nothing older.** Three pushed
|
||||
// while nothing is listening leave one on the stream, and it is the newest — the wire-level answer
|
||||
// to novox/hq issue 107, and the half of the drain that stops being the host's problem.
|
||||
func TestNatsANodeThatWasAwayGetsOnlyTheNewest(t *testing.T) {
|
||||
_, js := aBus(t)
|
||||
const node = "returning"
|
||||
_, private, _ := ed25519.GenerateKey(nil)
|
||||
|
||||
for _, id := range []string{"d1", "d2", "d3"} {
|
||||
if _, err := js.Publish(DeclareSubject(node),
|
||||
signedBy(t, private, []byte(`{"declared":"`+id+`"}`))); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
info, err := js.StreamInfo("NODES")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if info.State.Msgs != 1 {
|
||||
t.Fatalf("%d declarations survived for one node; a node that was away would apply a backlog "+
|
||||
"of things nobody wants any more", info.State.Msgs)
|
||||
}
|
||||
|
||||
theMeshMakes(t, js, node)
|
||||
feed := make(chan *nats.Msg, drainDepth)
|
||||
sub, err := js.ChanSubscribe(DeclareSubject(node), feed, nats.Bind("NODES", node))
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer func() { _ = sub.Unsubscribe() }()
|
||||
|
||||
select {
|
||||
case msg := <-feed:
|
||||
if declaredIn(msg.Data) != "d3" {
|
||||
t.Fatalf("the node was given %q rather than the newest", declaredIn(msg.Data))
|
||||
}
|
||||
case <-time.After(8 * time.Second):
|
||||
t.Fatal("the node that was away was given nothing")
|
||||
}
|
||||
}
|
||||
|
||||
// A report goes through the stream and a heartbeat does not: the one that must survive the
|
||||
// controller's store restarting is kept, and the one that must not is not.
|
||||
func TestNatsAReportIsKeptAndAHeartbeatIsNot(t *testing.T) {
|
||||
conn, js := aBus(t)
|
||||
bus := OverNATS{Conn: conn, JS: js}
|
||||
ctx := context.Background()
|
||||
|
||||
body, _ := json.Marshal(Report{Node: "anchor", Declared: "d1"})
|
||||
if err := bus.Report(ctx, "anchor", body); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
beat, _ := json.Marshal(Alive{Node: "anchor"})
|
||||
if err := bus.Alive(ctx, "anchor", beat); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
info, err := js.StreamInfo("CONTROL")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if info.State.Msgs != 1 {
|
||||
t.Fatalf("%d messages were kept; a report must be and a heartbeat must not", info.State.Msgs)
|
||||
}
|
||||
}
|
||||
@@ -3,33 +3,46 @@ package link
|
||||
import (
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
amqp "github.com/rabbitmq/amqp091-go"
|
||||
)
|
||||
|
||||
// said is one declaration as a test hands it over, with no transport under it — which is what the
|
||||
// seam bought: the drain's reasoning was reachable only through a real broker before.
|
||||
type said struct {
|
||||
body []byte
|
||||
handled bool
|
||||
}
|
||||
|
||||
func (s *said) Body() []byte { return s.body }
|
||||
func (s *said) Handled() error { s.handled = true; return nil }
|
||||
|
||||
func arriving(bodies ...string) chan Declaration {
|
||||
ch := make(chan Declaration, 8)
|
||||
for _, b := range bodies {
|
||||
ch <- &said{body: []byte(b)}
|
||||
}
|
||||
return ch
|
||||
}
|
||||
|
||||
// A machine asked to be five things becomes the last one: what is already waiting supersedes what
|
||||
// arrived first, and everything set aside is named so it can be reported.
|
||||
func TestWhatIsAlreadyWaitingSupersedesWhatArrivedFirst(t *testing.T) {
|
||||
deliveries := make(chan amqp.Delivery, 8)
|
||||
for _, id := range []string{"two", "three", "four"} {
|
||||
deliveries <- amqp.Delivery{Body: []byte(id)}
|
||||
waiting := arriving("two", "three", "four")
|
||||
apply, superseded := newest(waiting, &said{body: []byte("one")}, 50*time.Millisecond)
|
||||
if string(apply.Body()) != "four" {
|
||||
t.Fatalf("applied %q, not the newest", apply.Body())
|
||||
}
|
||||
apply, superseded := newest(deliveries, amqp.Delivery{Body: []byte("one")}, 50*time.Millisecond)
|
||||
if string(apply.Body) != "four" {
|
||||
t.Fatalf("applied %q, not the newest", apply.Body)
|
||||
}
|
||||
if len(superseded) != 3 || string(superseded[0].Body) != "one" || string(superseded[2].Body) != "three" {
|
||||
if len(superseded) != 3 || string(superseded[0].Body()) != "one" ||
|
||||
string(superseded[2].Body()) != "three" {
|
||||
t.Fatalf("set aside %d: %v", len(superseded), superseded)
|
||||
}
|
||||
}
|
||||
|
||||
// One declaration with nothing behind it is applied as it always was, after the window.
|
||||
func TestALoneDeclarationIsAppliedAfterTheWindow(t *testing.T) {
|
||||
deliveries := make(chan amqp.Delivery, 1)
|
||||
began := time.Now()
|
||||
apply, superseded := newest(deliveries, amqp.Delivery{Body: []byte("only")}, 30*time.Millisecond)
|
||||
if string(apply.Body) != "only" || len(superseded) != 0 {
|
||||
t.Fatalf("got %q with %d set aside", apply.Body, len(superseded))
|
||||
apply, superseded := newest(arriving(), &said{body: []byte("only")}, 30*time.Millisecond)
|
||||
if string(apply.Body()) != "only" || len(superseded) != 0 {
|
||||
t.Fatalf("got %q with %d set aside", apply.Body(), len(superseded))
|
||||
}
|
||||
if time.Since(began) < 30*time.Millisecond {
|
||||
t.Fatal("did not wait the window for a straggler")
|
||||
@@ -38,13 +51,13 @@ func TestALoneDeclarationIsAppliedAfterTheWindow(t *testing.T) {
|
||||
|
||||
// A straggler within the window is taken; one after it is the next push.
|
||||
func TestAStragglerWithinTheWindowIsTaken(t *testing.T) {
|
||||
deliveries := make(chan amqp.Delivery, 2)
|
||||
waiting := arriving()
|
||||
go func() {
|
||||
time.Sleep(20 * time.Millisecond)
|
||||
deliveries <- amqp.Delivery{Body: []byte("late")}
|
||||
waiting <- &said{body: []byte("late")}
|
||||
}()
|
||||
apply, superseded := newest(deliveries, amqp.Delivery{Body: []byte("first")}, 100*time.Millisecond)
|
||||
if string(apply.Body) != "late" || len(superseded) != 1 {
|
||||
t.Fatalf("got %q with %d set aside", apply.Body, len(superseded))
|
||||
apply, superseded := newest(waiting, &said{body: []byte("first")}, 100*time.Millisecond)
|
||||
if string(apply.Body()) != "late" || len(superseded) != 1 {
|
||||
t.Fatalf("got %q with %d set aside", apply.Body(), len(superseded))
|
||||
}
|
||||
}
|
||||
|
||||
+42
-108
@@ -6,10 +6,7 @@ import (
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"net/url"
|
||||
"time"
|
||||
|
||||
amqp "github.com/rabbitmq/amqp091-go"
|
||||
)
|
||||
|
||||
// ErrForged is what a node returns for a declaration whose signature is not the mesh's.
|
||||
@@ -33,6 +30,10 @@ type Membership struct {
|
||||
Fingerprint string
|
||||
Password string
|
||||
Signer ed25519.PublicKey
|
||||
// Transport is which bus this node speaks (hearing.go). Empty is the one the mesh runs on
|
||||
// today, which is every node until the rollout — so a membership recorded before any of this
|
||||
// existed reads as correct rather than as unset.
|
||||
Transport string
|
||||
}
|
||||
|
||||
// Applier is what the host does with a declaration that has been proved to come from the mesh.
|
||||
@@ -190,108 +191,62 @@ func Run(ctx context.Context, m Membership, apply Applier, say Announce, timeout
|
||||
if say == nil {
|
||||
say = func(string) {}
|
||||
}
|
||||
config, err := PinnedConfig(m.Fingerprint)
|
||||
link, err := Open(ctx, m, timeout)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer link.Close()
|
||||
|
||||
dsn := fmt.Sprintf("amqps://%s:%s@%s/",
|
||||
url.QueryEscape(m.Node), url.QueryEscape(m.Password), m.Broker)
|
||||
conn, err := amqp.DialConfig(dsn, amqp.Config{
|
||||
TLSClientConfig: config,
|
||||
Dial: amqp.DefaultDial(timeout),
|
||||
// Kept short so a node that has silently lost its route notices, rather than holding a
|
||||
// connection the broker forgot about and believing it is still in the mesh.
|
||||
Heartbeat: 10 * time.Second,
|
||||
})
|
||||
if err != nil {
|
||||
if errors.Is(err, ErrWrongCertificate) {
|
||||
return err
|
||||
}
|
||||
return fmt.Errorf("cannot reach the broker at %s: %w", m.Broker, err)
|
||||
}
|
||||
defer conn.Close()
|
||||
|
||||
channel, err := conn.Channel()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer channel.Close()
|
||||
|
||||
queue := QueueFor(m.Node)
|
||||
if _, err := channel.QueueDeclare(queue, true, false, false, false, nil); err != nil {
|
||||
return fmt.Errorf("cannot declare this node's queue %s: %w", queue, err)
|
||||
}
|
||||
|
||||
// Applying is one at a time — two at once would race on the same filesystem — but SEEING is
|
||||
// not: with a prefetch of one the host could never know that a newer declaration was already
|
||||
// waiting, and so applied every one of a backlog in turn, at the better part of a minute each,
|
||||
// becoming things nobody wanted any more (novox/hq issue 031). A window of unacknowledged
|
||||
// deliveries lets it drain to the newest; each declaration still survives a restart on the
|
||||
// broker until it is acknowledged, which happens only after it is applied or set aside.
|
||||
if err := channel.Qos(drainDepth, 0, false); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
deliveries, err := channel.ConsumeWithContext(ctx, queue, "", false, false, false, false, nil)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
// Said, because it is the event anybody watching actually wants. Without it a node logs
|
||||
// every failure and nothing on success, so a log full of "trying again" and then silence
|
||||
// reads as still broken when it means the opposite.
|
||||
say("in the mesh, consuming " + queue)
|
||||
// Said, because it is the event anybody watching actually wants. Without it a node logs every
|
||||
// failure and nothing on success, so a log full of "trying again" and then silence reads as
|
||||
// still broken when it means the opposite.
|
||||
say("in the mesh, hearing what this node should be")
|
||||
|
||||
// A word every so often, so the mesh can tell a node that is quiet from one that is gone.
|
||||
// Cheap on purpose: it carries a name and nothing else, because anything more would be a
|
||||
// report, and reports are rare where this is constant.
|
||||
beat := time.NewTicker(AliveEvery)
|
||||
defer beat.Stop()
|
||||
publishAlive(ctx, OverCurrent{Channel: channel}, m, say, timeout)
|
||||
publishAlive(ctx, link, m, say, timeout)
|
||||
|
||||
closed := conn.NotifyClose(make(chan *amqp.Error, 1))
|
||||
|
||||
// Published mandatory, so the broker hands back anything it cannot route rather than
|
||||
// dropping it. Without this a report goes to an exchange with no matching binding, the
|
||||
// publisher is told nothing, and the mesh believes this node never answered while the node
|
||||
// believes it did — which is what happened when `report` was left unbound on the other side.
|
||||
returned := channel.NotifyReturn(make(chan amqp.Return, 4))
|
||||
go func() {
|
||||
for r := range returned {
|
||||
say(fmt.Sprintf("the broker could not route this node's %s: %s (%d %s)",
|
||||
r.RoutingKey, r.Exchange, r.ReplyCode, r.ReplyText))
|
||||
}
|
||||
}()
|
||||
declarations := link.Declarations()
|
||||
|
||||
for {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return nil
|
||||
case <-beat.C:
|
||||
publishAlive(ctx, OverCurrent{Channel: channel}, m, say, timeout)
|
||||
publishAlive(ctx, link, m, say, timeout)
|
||||
case unasked := <-outbox:
|
||||
// Said without having been asked: a reconcile found what an adopted node holds, or
|
||||
// its firewall, changed since it last said.
|
||||
published := publishReport(ctx, OverCurrent{Channel: channel}, m, unasked.Report, say, timeout)
|
||||
published := publishReport(ctx, link, m, unasked.Report, say, timeout)
|
||||
if unasked.Done != nil {
|
||||
unasked.Done(published)
|
||||
}
|
||||
case reason := <-closed:
|
||||
return fmt.Errorf("the link closed: %v", reason)
|
||||
case delivery, ok := <-deliveries:
|
||||
case reason := <-link.Lost():
|
||||
return reason
|
||||
case declaration, ok := <-declarations:
|
||||
if !ok {
|
||||
return errors.New("the broker stopped delivering")
|
||||
// The link's own reason, when it has managed to say one: "stopped delivering" on
|
||||
// its own says nothing about why, and why is the whole of what an operator wants.
|
||||
select {
|
||||
case reason := <-link.Lost():
|
||||
return reason
|
||||
default:
|
||||
return errors.New("the mesh stopped sending this node declarations")
|
||||
}
|
||||
}
|
||||
// Whatever else is already waiting supersedes this one. Each set-aside declaration
|
||||
// is reported as such, then acknowledged unapplied.
|
||||
delivery, superseded := newest(deliveries, delivery, drainWindow)
|
||||
// Whatever else is already waiting supersedes this one. Each set-aside declaration is
|
||||
// reported as such, then settled unapplied.
|
||||
declaration, superseded := newest(declarations, declaration, drainWindow)
|
||||
for _, old := range superseded {
|
||||
say("set aside a declaration: a newer one arrived with it")
|
||||
publishReport(ctx, OverCurrent{Channel: channel}, m, Report{Node: m.Node, Declared: declaredIn(old.Body),
|
||||
Superseded: declaredIn(delivery.Body)}, say, timeout)
|
||||
_ = old.Ack(false)
|
||||
publishReport(ctx, link, m, Report{Node: m.Node, Declared: declaredIn(old.Body()),
|
||||
Superseded: declaredIn(declaration.Body())}, say, timeout)
|
||||
_ = old.Handled()
|
||||
}
|
||||
report := handle(ctx, m, apply, delivery)
|
||||
report := handleBody(ctx, m, declaration.Body(), apply)
|
||||
switch {
|
||||
case report.Refused != "":
|
||||
say("refused a declaration: " + report.Refused)
|
||||
@@ -300,12 +255,12 @@ func Run(ctx context.Context, m Membership, apply Applier, say Announce, timeout
|
||||
default:
|
||||
say(fmt.Sprintf("applied %d resource(s)", len(report.Applied)))
|
||||
}
|
||||
publishReport(ctx, OverCurrent{Channel: channel}, m, report, say, timeout)
|
||||
// Acknowledged after the report is published. A node that dies between applying and
|
||||
// reporting leaves the declaration on the broker and applies it again on return,
|
||||
publishReport(ctx, link, m, report, say, timeout)
|
||||
// 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.
|
||||
_ = delivery.Ack(false)
|
||||
_ = declaration.Handled()
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -333,12 +288,12 @@ const (
|
||||
// three deliveries, whatever the stream later retains. So this is narrowed at the rollout, not
|
||||
// deleted — and saying which half goes is worth more than a note that it "can probably be
|
||||
// removed", which is how a load-bearing window gets deleted by somebody in a hurry.
|
||||
func newest(deliveries <-chan amqp.Delivery, first amqp.Delivery, window time.Duration) (amqp.Delivery, []amqp.Delivery) {
|
||||
func newest(arriving <-chan Declaration, first Declaration, window time.Duration) (Declaration, []Declaration) {
|
||||
latest := first
|
||||
var superseded []amqp.Delivery
|
||||
var superseded []Declaration
|
||||
for {
|
||||
select {
|
||||
case next, ok := <-deliveries:
|
||||
case next, ok := <-arriving:
|
||||
if !ok {
|
||||
return latest, superseded
|
||||
}
|
||||
@@ -367,10 +322,6 @@ func declaredIn(body []byte) string {
|
||||
return d.Declared
|
||||
}
|
||||
|
||||
func handle(ctx context.Context, m Membership, apply Applier, delivery amqp.Delivery) Report {
|
||||
return handleBody(ctx, m, delivery.Body, apply)
|
||||
}
|
||||
|
||||
// handleBody is the whole of deciding whether to trust a message, separated from the broker so it
|
||||
// can be tested as the security check it is rather than as message plumbing.
|
||||
func handleBody(ctx context.Context, m Membership, body []byte, apply Applier) Report {
|
||||
@@ -393,30 +344,13 @@ func handleBody(ctx context.Context, m Membership, body []byte, apply Applier) R
|
||||
// report a command makes rather than the running host — a rekey (novox/hq ADR 0105). The same
|
||||
// account, the same pinned certificate and the same exchange as the running host's reports.
|
||||
func Publish(ctx context.Context, m Membership, report Report, timeout time.Duration) error {
|
||||
config, err := PinnedConfig(m.Fingerprint)
|
||||
link, err := Open(ctx, m, timeout)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
dsn := fmt.Sprintf("amqps://%s:%s@%s/",
|
||||
url.QueryEscape(m.Node), url.QueryEscape(m.Password), m.Broker)
|
||||
conn, err := amqp.DialConfig(dsn, amqp.Config{
|
||||
TLSClientConfig: config,
|
||||
Dial: amqp.DefaultDial(timeout),
|
||||
})
|
||||
if err != nil {
|
||||
if errors.Is(err, ErrWrongCertificate) {
|
||||
return err
|
||||
}
|
||||
return fmt.Errorf("cannot reach the broker at %s: %w", m.Broker, err)
|
||||
}
|
||||
defer conn.Close()
|
||||
channel, err := conn.Channel()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer channel.Close()
|
||||
defer link.Close()
|
||||
var said string
|
||||
if !publishReport(ctx, OverCurrent{Channel: channel}, m, report, func(s string) { said = s }, timeout) {
|
||||
if !publishReport(ctx, link, m, report, func(s string) { said = s }, timeout) {
|
||||
return errors.New(said)
|
||||
}
|
||||
return nil
|
||||
|
||||
Reference in New Issue
Block a user