Compare commits
5
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
22845a5296 | ||
|
|
e7da39de57 | ||
|
|
e6ddc59cde | ||
|
|
6c5dfd0c25 | ||
|
|
775df79893 |
@@ -346,9 +346,16 @@ func whoResolves(ctx context.Context, open *stores, requirement string) (
|
||||
refused := map[string]string{}
|
||||
for _, n := range nodes {
|
||||
plan, _, err := planFor(ctx, open, n.Name)
|
||||
if err != nil {
|
||||
switch {
|
||||
case unresolvable(err):
|
||||
refused[n.Name] = err.Error()
|
||||
continue
|
||||
case err != nil:
|
||||
// Not a node that does not resolve — a question that went unanswered. Recording it as a
|
||||
// refusal would take the machine off the private network, and the generator that reads
|
||||
// this would then write a roster and a filter without it (novox/hq 04-ISSUES/152).
|
||||
return nil, nil, fmt.Errorf("whether %s answers %q cannot be read: %w",
|
||||
n.Name, requirement, err)
|
||||
}
|
||||
for _, m := range plan.Modules {
|
||||
for _, offered := range m.Offers() {
|
||||
|
||||
@@ -25,7 +25,36 @@ import (
|
||||
// cheapest next step. That is how novox/hq ADR 0001 records `hal/sdk` reaching 34,636:
|
||||
// nothing in it was wrong, and no one edit was the one that should have been a new file.
|
||||
|
||||
// notResolvable marks the one failure in planFor that is a statement about the node: its assigned
|
||||
// modules do not compose. Every other failure means the mesh could not be *asked* — the store was
|
||||
// unreachable, a key could not be read — and says nothing about the node at all.
|
||||
//
|
||||
// The distinction exists because three callers gather something across every machine and must carry
|
||||
// on when one machine's set is broken. Each of them read a plain error as "their set does not
|
||||
// resolve", and so read a store that was briefly unreachable as a machine that runs nothing. On the
|
||||
// roster of routed names that is not a degraded answer but a false one: it states, to every machine
|
||||
// at once, that another machine's names do not exist. A control node spent hours replacing every
|
||||
// container it ran, on a six-minute cycle, because each pass restarted the store this is read from,
|
||||
// the read failed, one name left the roster, and the roster is part of every container's identity
|
||||
// (novox/hq 04-ISSUES/152, and 04-ISSUES/151 for why a changed roster is a changed container).
|
||||
//
|
||||
// So: skip a node that cannot resolve, and never a node that could not be read.
|
||||
type notResolvable struct{ err error }
|
||||
|
||||
func (n notResolvable) Error() string { return n.err.Error() }
|
||||
func (n notResolvable) Unwrap() error { return n.err }
|
||||
|
||||
// unresolvable reports whether err is a node's own set failing to compose, rather than the mesh
|
||||
// being unable to answer.
|
||||
func unresolvable(err error) bool {
|
||||
var n notResolvable
|
||||
return errors.As(err, &n)
|
||||
}
|
||||
|
||||
// planFor works out everything a node should run, from what was assigned to it.
|
||||
//
|
||||
// A failure to compose the node's own modules is wrapped as notResolvable; every other failure is
|
||||
// returned as it is. Callers gathering across the mesh must tell them apart — see notResolvable.
|
||||
func planFor(ctx context.Context, open *stores, nodeName string) (catalogue.Resolution, catalogue.SettingsBy, error) {
|
||||
inv := open.inventory
|
||||
shelf, err := inv.Catalogue(ctx)
|
||||
@@ -95,7 +124,9 @@ func planFor(ctx context.Context, open *stores, nodeName string) (catalogue.Reso
|
||||
At: onNetwork[nodeName], PublicDomain: publicDomain,
|
||||
Account: who.Account, AccountHome: who.AccountHome}, world)
|
||||
if err != nil {
|
||||
return catalogue.Resolution{}, nil, err
|
||||
// The node's own set does not compose. Marked, because this is the only failure here that
|
||||
// a mesh-wide gatherer may pass over — see notResolvable.
|
||||
return catalogue.Resolution{}, nil, notResolvable{err}
|
||||
}
|
||||
|
||||
// The credential for each thing this node takes from elsewhere. Made once and kept, so the
|
||||
@@ -163,8 +194,13 @@ func planFor(ctx context.Context, open *stores, nodeName string) (catalogue.Reso
|
||||
if len(stray) > 0 {
|
||||
// Somebody set something that reaches no file. Said here rather than discovered by the
|
||||
// machine not behaving differently, which is the slowest way there is.
|
||||
return catalogue.Resolution{}, nil, fmt.Errorf(
|
||||
"these settings reach nothing:\n - %s", strings.Join(stray, "\n - "))
|
||||
//
|
||||
// Marked like a set that will not compose, and for the same reason: it is a standing fact
|
||||
// about this node's own configuration, not a question the mesh could not answer. A gatherer
|
||||
// passes over it as it always did — one node's stray setting must not stop every other node
|
||||
// being described (novox/hq 04-ISSUES/152).
|
||||
return catalogue.Resolution{}, nil, notResolvable{fmt.Errorf(
|
||||
"these settings reach nothing:\n - %s", strings.Join(stray, "\n - "))}
|
||||
}
|
||||
return resolved, settings, nil
|
||||
}
|
||||
@@ -671,11 +707,17 @@ func renderingFor(ctx context.Context, open *stores, node string,
|
||||
// routed name only because it carried a label the mesh composed, never because the mesh knows what
|
||||
// "route" means. A node that does not resolve is skipped, so one machine's broken set does not cost
|
||||
// the rest their names.
|
||||
//
|
||||
// **A node that could not be READ is a different matter and is raised.** Skipping one states, to
|
||||
// every machine at once, that its names do not exist — and since the roster is part of every
|
||||
// container's identity, that withdraws them and replaces every container (novox/hq 04-ISSUES/152,
|
||||
// 151). So every failure here says which machine and which read, because the alternative is a
|
||||
// mesh-wide refusal with nothing named in it.
|
||||
func routeNamesInTheMesh(ctx context.Context, open *stores) (map[string]string, error) {
|
||||
inv := open.inventory
|
||||
places, err := inv.Overlays(ctx)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
return nil, fmt.Errorf("where the machines are cannot be read: %w", err)
|
||||
}
|
||||
address := map[string]string{}
|
||||
for _, p := range places {
|
||||
@@ -686,14 +728,22 @@ func routeNamesInTheMesh(ctx context.Context, open *stores) (map[string]string,
|
||||
|
||||
nodes, err := inv.Nodes(ctx)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
return nil, fmt.Errorf("which machines the mesh has cannot be read: %w", err)
|
||||
}
|
||||
|
||||
out := map[string]string{}
|
||||
for _, n := range nodes {
|
||||
plan, settings, err := planFor(ctx, open, n.Name)
|
||||
if err != nil {
|
||||
switch {
|
||||
case unresolvable(err):
|
||||
// Their set does not compose, so they serve no names. Passed over, so one machine's
|
||||
// broken set does not cost the rest theirs.
|
||||
continue
|
||||
case err != nil:
|
||||
// The mesh could not be asked. Returning the roster without this machine's names would
|
||||
// state that they do not exist — to every machine, and indistinguishably from the
|
||||
// operator having withdrawn them (novox/hq 04-ISSUES/152).
|
||||
return nil, fmt.Errorf("the names %s serves cannot be read: %w", n.Name, err)
|
||||
}
|
||||
for _, m := range plan.Modules {
|
||||
for to := range m.Contributes {
|
||||
@@ -823,11 +873,17 @@ func grantsFor(ctx context.Context, open *stores, node string) ([]catalogue.Gran
|
||||
out := make([]catalogue.Grant, 0, len(issued))
|
||||
for _, s := range issued {
|
||||
plan, settings, err := planFor(ctx, open, s.Consumer)
|
||||
if err != nil {
|
||||
switch {
|
||||
case unresolvable(err):
|
||||
// Their set does not resolve. Skipped rather than fatal: this node is not the place
|
||||
// to report another machine's problem, and a grant for something that is not going to
|
||||
// run would have the provider create a user nothing uses.
|
||||
continue
|
||||
case err != nil:
|
||||
// The mesh could not be asked what they wanted, which is not the same as their wanting
|
||||
// nothing — and withholding a grant on that reading takes a consumer's access away
|
||||
// (novox/hq 04-ISSUES/152).
|
||||
return nil, fmt.Errorf("what %s asked of %s cannot be read: %w", s.Consumer, s.Name, err)
|
||||
}
|
||||
values, asks, err := plan.ContributionsFrom(s.Name, s.ConsumerModule, settings)
|
||||
if err != nil {
|
||||
|
||||
@@ -717,6 +717,12 @@ func raiseTheBus(ctx context.Context, inv *inventory.Inventory, address string)
|
||||
broker.BareAddress(address), err)
|
||||
}
|
||||
defer js.Close()
|
||||
// What the raise decided not to fail over. Said, for the reason everything else here is said:
|
||||
// a consumer kept as it was is a difference between what the mesh asked for and what the bus
|
||||
// holds, and one nobody would find by reading either (novox/hq 04-ISSUES/156).
|
||||
js.Note = func(format string, args ...any) {
|
||||
fmt.Printf(" "+format+"\n", args...)
|
||||
}
|
||||
|
||||
// **Its own user, before anything else.** The controller's account is created by the installer at
|
||||
// a bootstrap password, before there is a controller to mint one — so nothing recorded a hash for
|
||||
|
||||
@@ -0,0 +1,99 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"strings"
|
||||
"testing"
|
||||
)
|
||||
|
||||
// A node's own set failing to compose, and the mesh being unable to answer at all, are different
|
||||
// things, and only the first may be passed over when something is gathered across every machine
|
||||
// (novox/hq 04-ISSUES/152). These pin that distinction where the three gatherers rely on it.
|
||||
|
||||
func TestASetThatDoesNotComposeIsMarkedAsTheNodesOwnProblem(t *testing.T) {
|
||||
open := aMesh(t)
|
||||
one, two := rivals()
|
||||
register(t, open, one)
|
||||
register(t, open, two)
|
||||
for _, m := range []string{one.Module, two.Module} {
|
||||
if _, err := open.inventory.Assign(t.Context(), "laptop", m); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
|
||||
_, _, err := planFor(t.Context(), open, "laptop")
|
||||
if err == nil {
|
||||
t.Fatal("two modules claiming one seat composed anyway")
|
||||
}
|
||||
if !unresolvable(err) {
|
||||
t.Fatalf("a set that cannot compose was not marked as the node's own problem: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestAStoreThatCannotBeReadIsNotANodeThatDoesNotCompose(t *testing.T) {
|
||||
open := aMesh(t)
|
||||
|
||||
// Nothing is wrong with anchor. The question simply cannot be asked.
|
||||
stopped, cancel := context.WithCancel(t.Context())
|
||||
cancel()
|
||||
|
||||
_, _, err := planFor(stopped, open, "anchor")
|
||||
if err == nil {
|
||||
t.Fatal("a plan composed against a store that could not be read")
|
||||
}
|
||||
if unresolvable(err) {
|
||||
t.Fatalf("a question the mesh could not answer was read as a node that runs nothing: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestOneIncoherentNodeDoesNotCostTheRestTheirNames(t *testing.T) {
|
||||
open := aMesh(t)
|
||||
one, two := rivals()
|
||||
register(t, open, one)
|
||||
register(t, open, two)
|
||||
for _, m := range []string{one.Module, two.Module} {
|
||||
if _, err := open.inventory.Assign(t.Context(), "laptop", m); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
|
||||
// laptop cannot compose. That is laptop's problem and nobody else's: the roster is still
|
||||
// answerable, and anchor keeps whatever it serves.
|
||||
if _, err := routeNamesInTheMesh(t.Context(), open); err != nil {
|
||||
t.Fatalf("one node's broken set cost the whole mesh its roster: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestARosterIsNeverReturnedWithNamesItCouldNotRead(t *testing.T) {
|
||||
open := aMesh(t)
|
||||
|
||||
stopped, cancel := context.WithCancel(t.Context())
|
||||
cancel()
|
||||
|
||||
names, err := routeNamesInTheMesh(stopped, open)
|
||||
if err == nil {
|
||||
t.Fatalf("a roster was composed from a store that could not be read: %v", names)
|
||||
}
|
||||
// The failure must be raised, not turned into an absence. A roster missing a machine's names
|
||||
// is indistinguishable, on every machine that receives it, from the operator withdrawing them —
|
||||
// and because the roster is part of every container's identity, it replaces all of them.
|
||||
if names != nil {
|
||||
t.Fatalf("a partial roster was returned beside the error: %v", names)
|
||||
}
|
||||
}
|
||||
|
||||
// Kept so the reason survives the next person reading it: the message the gatherer raises must say
|
||||
// which machine could not be read, or the operator is left with a mesh-wide failure and no name.
|
||||
func TestTheRaisedFailureNamesTheMachineItCouldNotRead(t *testing.T) {
|
||||
open := aMesh(t)
|
||||
stopped, cancel := context.WithCancel(t.Context())
|
||||
cancel()
|
||||
|
||||
_, err := routeNamesInTheMesh(stopped, open)
|
||||
if err == nil {
|
||||
t.Fatal("no failure was raised")
|
||||
}
|
||||
if !strings.Contains(err.Error(), "cannot be read") {
|
||||
t.Fatalf("the failure does not say the mesh could not be read: %v", err)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,178 @@
|
||||
package broker
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"os"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/nats-io/nats.go"
|
||||
)
|
||||
|
||||
const twoSeconds = 2 * time.Second
|
||||
|
||||
// A running mesh already holds consumers made before the delivery subject carried the stream
|
||||
// (novox/hq 04-ISSUES/146). The server will not change a push consumer's delivery subject in place,
|
||||
// so bringing one to match must replace it — and must not replay what it already acknowledged
|
||||
// (novox/hq 04-ISSUES/156).
|
||||
//
|
||||
// docker run -d --rm --name t -p 14231:4222 nats:2.10-alpine -js
|
||||
// MESH_TEST_NATS=nats://127.0.0.1:14231 go test ./internal/broker/ -run TestUpgrading
|
||||
func TestUpgradingAConsumerWhoseDeliverySubjectMoved(t *testing.T) {
|
||||
url := os.Getenv("MESH_TEST_NATS")
|
||||
if url == "" {
|
||||
t.Skip("MESH_TEST_NATS unset")
|
||||
}
|
||||
js, err := Dial(url)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer js.Close()
|
||||
|
||||
// The stream exactly as the mesh's own is — one declaration per node, always the newest.
|
||||
// Reproduced rather than approximated: the first version of this test used a plain stream and
|
||||
// a plain consumer, and the server accepted the update it refuses in a running mesh, so the
|
||||
// test passed against the very code that was crash-looping on the control node.
|
||||
const stream, name = "NODES", "novox"
|
||||
subject := "mesh.node." + name + ".declare"
|
||||
_ = js.js.DeleteStream(stream)
|
||||
if _, err := js.js.AddStream(&nats.StreamConfig{
|
||||
Name: stream, Subjects: []string{"mesh.node.*.declare"},
|
||||
MaxMsgsPerSubject: 1, Storage: nats.MemoryStorage,
|
||||
}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer func() { _ = js.js.DeleteStream(stream) }()
|
||||
|
||||
for i := 0; i < 6; i++ {
|
||||
if _, err := js.js.Publish(subject, []byte(fmt.Sprint(i))); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
|
||||
// The consumer as a running mesh holds it: made before the subject carried the stream, and
|
||||
// otherwise exactly what NodeConsumer asks for.
|
||||
if _, err := js.js.AddConsumer(stream, &nats.ConsumerConfig{
|
||||
Durable: name, AckPolicy: nats.AckExplicitPolicy,
|
||||
AckWait: 300 * time.Second, MaxDeliver: -1,
|
||||
FilterSubject: subject,
|
||||
DeliverSubject: "_DELIVER." + name,
|
||||
}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
// It acknowledged the first four. Those must not come back.
|
||||
sub, err := js.js.SubscribeSync(subject, nats.Bind(stream, name))
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
for i := 0; i < 1; i++ {
|
||||
m, err := sub.NextMsg(twoSeconds)
|
||||
if err != nil {
|
||||
t.Fatalf("message %d never arrived: %v", i, err)
|
||||
}
|
||||
if err := m.AckSync(); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
// **The subscription stays up.** In a running mesh the machine is attached to this consumer
|
||||
// the whole time — that is what a node listening for its declaration IS. The first version of
|
||||
// this test unsubscribed first, and the server then accepted an update it refuses while a
|
||||
// subscriber is bound, so the test passed against the code that was crash-looping.
|
||||
defer func() { _ = sub.Unsubscribe() }()
|
||||
|
||||
// Now the upgrade: the consumer the controller asserts on every start, with the subject that
|
||||
// carries the stream.
|
||||
want := NodeConsumer(name)
|
||||
var notes []string
|
||||
js.Note = func(f string, a ...any) { notes = append(notes, fmt.Sprintf(f, a...)) }
|
||||
|
||||
if err := js.EnsureConsumer(want); err != nil {
|
||||
t.Fatalf("a consumer the mesh already held could not be brought to match, which is the "+
|
||||
"control plane failing to start: %v", err)
|
||||
}
|
||||
|
||||
info, err := js.js.ConsumerInfo(stream, name)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
// It KEEPS the subject it had. Moving it would need the holder's grant to have widened first,
|
||||
// and that grant travels in the bus's user list, which a machine applies minutes later.
|
||||
if got := info.Config.DeliverSubject; got != "_DELIVER."+name {
|
||||
t.Fatalf("the consumer a machine is bound to was moved to %q; a machine not yet allowed "+
|
||||
"to subscribe there is a machine that hears nothing", got)
|
||||
}
|
||||
if len(notes) != 1 {
|
||||
t.Fatalf("keeping it was not reported, so it would be invisible: %v", notes)
|
||||
}
|
||||
if !strings.Contains(notes[0], "keeps working") {
|
||||
t.Fatalf("the note does not say the consumer still works: %q", notes[0])
|
||||
}
|
||||
|
||||
// And the machine bound to it is still being delivered to — the point of keeping it.
|
||||
if _, err := js.js.Publish(subject, []byte("after the assertion")); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
m, err := sub.NextMsg(twoSeconds)
|
||||
if err != nil {
|
||||
t.Fatalf("the machine stopped hearing its declarations after the assertion: %v", err)
|
||||
}
|
||||
if string(m.Data) != "after the assertion" {
|
||||
t.Fatalf("delivered %q", m.Data)
|
||||
}
|
||||
|
||||
// Asserting again is a no-op, or the controller crash-loops on its own restart.
|
||||
if err := js.EnsureConsumer(want); err != nil {
|
||||
t.Fatalf("the second assertion failed: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
// And where nothing is bound, the subject DOES move — that is 04-ISSUES/146's fix, which this must
|
||||
// not undo. The controller's own two consumers are in exactly this position: it asserts them before
|
||||
// it subscribes.
|
||||
func TestAConsumerNothingIsBoundToDoesMove(t *testing.T) {
|
||||
url := os.Getenv("MESH_TEST_NATS")
|
||||
if url == "" {
|
||||
t.Skip("MESH_TEST_NATS unset")
|
||||
}
|
||||
js, err := Dial(url)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer js.Close()
|
||||
|
||||
const stream, name = "NODES", "shanks"
|
||||
subject := "mesh.node." + name + ".declare"
|
||||
_ = js.js.DeleteStream(stream)
|
||||
if _, err := js.js.AddStream(&nats.StreamConfig{
|
||||
Name: stream, Subjects: []string{"mesh.node.*.declare"},
|
||||
MaxMsgsPerSubject: 1, Storage: nats.MemoryStorage,
|
||||
}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer func() { _ = js.js.DeleteStream(stream) }()
|
||||
|
||||
if _, err := js.js.AddConsumer(stream, &nats.ConsumerConfig{
|
||||
Durable: name, AckPolicy: nats.AckExplicitPolicy,
|
||||
AckWait: 300 * time.Second, MaxDeliver: -1,
|
||||
FilterSubject: subject,
|
||||
DeliverSubject: "_DELIVER." + name,
|
||||
}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
want := NodeConsumer(name)
|
||||
if err := js.EnsureConsumer(want); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
info, err := js.js.ConsumerInfo(stream, name)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if got := info.Config.DeliverSubject; got != DeliverSubjectFor(want) {
|
||||
t.Fatalf("delivery subject is %q, wanted %q -- issue 146's fix no longer applies to a "+
|
||||
"consumer nothing is holding", got, DeliverSubjectFor(want))
|
||||
}
|
||||
}
|
||||
@@ -26,6 +26,16 @@ import (
|
||||
type JetStream struct {
|
||||
conn *nats.Conn
|
||||
js nats.JetStreamContext
|
||||
// Note is how this says something it decided not to fail over. Nil is silent, which is only
|
||||
// right for a caller that has no way to report; the controller sets it.
|
||||
Note func(string, ...any)
|
||||
}
|
||||
|
||||
// note reports without requiring a caller to have set one.
|
||||
func (j *JetStream) note(format string, args ...any) {
|
||||
if j.Note != nil {
|
||||
j.Note(format, args...)
|
||||
}
|
||||
}
|
||||
|
||||
// Dial connects and returns the controller's JetStream handle.
|
||||
@@ -200,9 +210,44 @@ func (j *JetStream) EnsureConsumer(c Consumer) error {
|
||||
want.DeliverSubject = DeliverSubjectFor(c)
|
||||
}
|
||||
|
||||
switch _, err := j.js.ConsumerInfo(c.Stream, c.Name); {
|
||||
switch have, err := j.js.ConsumerInfo(c.Stream, c.Name); {
|
||||
case err == nil:
|
||||
// Where an existing consumer starts is its history, not something an assertion may move:
|
||||
// the server refuses a changed deliver policy outright. Carried across, so asserting twice
|
||||
// is the no-op a restart depends on.
|
||||
want.DeliverPolicy = have.Config.DeliverPolicy
|
||||
want.OptStartSeq = have.Config.OptStartSeq
|
||||
want.OptStartTime = have.Config.OptStartTime
|
||||
|
||||
if _, err := j.js.UpdateConsumer(c.Stream, want); err != nil {
|
||||
// **A consumer that works is not replaced to make its name tidier**
|
||||
// (novox/hq 04-ISSUES/156).
|
||||
//
|
||||
// The server will not move a push consumer's delivery subject while a subscriber is
|
||||
// bound to it, and answers `consumer name already in use` — a message about the name,
|
||||
// for a conflict about the subject. A node is bound to its declaration consumer the
|
||||
// whole time it is up; that IS a node listening. So when 04-ISSUES/146 put the stream
|
||||
// into the subject, every node consumer in a running mesh became one this could not
|
||||
// bring to match, and the control plane crash-looped on the assertion it makes before
|
||||
// it serves. A fresh mesh showed nothing: nothing was bound.
|
||||
//
|
||||
// Kept rather than deleted and re-made. Re-making moves the subject, and a holder may
|
||||
// not be allowed to subscribe to the new one yet — the wider grant travels in the bus's
|
||||
// user list, which this same control plane composes and a machine applies minutes
|
||||
// later. Re-making here would have silenced every machine in the mesh, which is worse
|
||||
// than the collision it was fixing and harder to undo.
|
||||
//
|
||||
// Kept rather than fatal, which is what 146's change intended and did not do: the bare
|
||||
// subject it replaces still delivers, and it collides only where one holder has two
|
||||
// consumers of one name. That is the controller's own pair, and the controller is not
|
||||
// bound to them while it asserts, so those do move. A node has one consumer and nothing
|
||||
// to collide with.
|
||||
if have.Config.DeliverSubject != want.DeliverSubject {
|
||||
j.note("consumer %s on %s still delivers to %q and not %q: %v. It keeps working; "+
|
||||
"the subject moves on an assertion made while nothing is bound to it",
|
||||
c.Name, c.Stream, have.Config.DeliverSubject, want.DeliverSubject, err)
|
||||
return nil
|
||||
}
|
||||
return fmt.Errorf("bringing consumer %s on %s to match: %w", c.Name, c.Stream, err)
|
||||
}
|
||||
return nil
|
||||
|
||||
@@ -70,6 +70,27 @@ func knownFor(m Manifest, needs []Needed, node string) map[string]map[string]str
|
||||
return out
|
||||
}
|
||||
|
||||
// withOwnNames adds a module's own composed names to what it may name from one binding:
|
||||
// `${bound:<provision>:name}` and `:internal-name`, and for several contributions to one requirement
|
||||
// `:name-<local>` / `:internal-name-<local>`. Set over anything the provider serves under those keys:
|
||||
// what the module is called is the mesh's statement, not the provider's.
|
||||
func withOwnNames(values map[string]string, own map[string]any) {
|
||||
for _, key := range []string{"name", "internal-name"} {
|
||||
if v, ok := own[key].(string); ok {
|
||||
values[key] = v
|
||||
}
|
||||
}
|
||||
many, _ := own["names"].(map[string]any)
|
||||
for local, raw := range many {
|
||||
names, _ := raw.(map[string]any)
|
||||
for _, key := range []string{"name", "internal-name"} {
|
||||
if v, ok := names[key].(string); ok {
|
||||
values[key+"-"+local] = v
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// plainly renders a served value as a program would expect to read it.
|
||||
func plainly(value any) string {
|
||||
switch v := value.(type) {
|
||||
|
||||
@@ -574,7 +574,11 @@ func (r Resolution) compose(with Rendering, owner map[string]string) ([]map[stri
|
||||
}
|
||||
found = here
|
||||
}
|
||||
file, err := boundFile(*found, m.Binds[to], ConsumerIdentity(r.Node, IdentitySource(m.Slug, m.Module)))
|
||||
own, err := r.ownNames(m, to, with.Settings[m.Module])
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
file, err := boundFile(*found, m.Binds[to], ConsumerIdentity(r.Node, IdentitySource(m.Slug, m.Module)), own)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
@@ -637,6 +641,35 @@ func (r Resolution) compose(with Rendering, owner map[string]string) ([]map[stri
|
||||
}
|
||||
// And what its bindings say, for the half of a connection that is not secret.
|
||||
known := knownFor(m, r.Needs, r.Node)
|
||||
// A requirement answered on this same machine is not in r.Needs — its binding file is
|
||||
// written from `here` (above) — and so `${bound:…}` could not name it, though the file
|
||||
// beside it said the same facts. Filled from the same answer, so the two cannot disagree.
|
||||
for _, want := range m.Wants() {
|
||||
if _, has := known[want]; has {
|
||||
continue
|
||||
}
|
||||
answered, err := here(r, want, with)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if answered == nil {
|
||||
continue
|
||||
}
|
||||
local := *answered
|
||||
local.For = m.Module
|
||||
for provision, values := range knownFor(m, []Needed{local}, r.Node) {
|
||||
known[provision] = values
|
||||
}
|
||||
}
|
||||
// And what the module is called through each requirement it contributes to (novox/hq
|
||||
// 04-ISSUES/122) — the same composition its binding file carries.
|
||||
for provision, values := range known {
|
||||
own, err := r.ownNames(m, provision, with.Settings[m.Module])
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
withOwnNames(values, own)
|
||||
}
|
||||
// And where this node places the directories the module declared without a path
|
||||
// (novox/hq ADR 0112) — resolved once per module, named by ${dir:…} from any resource.
|
||||
dirs := dirsFor(m, with)
|
||||
@@ -1089,21 +1122,11 @@ func (r Resolution) contributions(settings SettingsBy, grants []Grant,
|
||||
// Settings reach a contribution the same way they reach a file. A route's hostname is
|
||||
// exactly the kind of thing that differs between one mesh and the next, and a module
|
||||
// that could not have it set would have to be edited to be reused.
|
||||
values, err := settle(m.Contributes[to], settings[m.Module], nil,
|
||||
values, err := r.composed(m, m.Contributes[to], settings[m.Module],
|
||||
m.Module+" contributing to "+to)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("%s contributing to %s: %w", m.Module, to, err)
|
||||
return nil, err
|
||||
}
|
||||
reaches, err := Reaches(m, settings[m.Module])
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("%s contributing to %s: %w", m.Module, to, err)
|
||||
}
|
||||
blocks, err := Endpoints(m, settings[m.Module])
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("%s contributing to %s: %w", m.Module, to, err)
|
||||
}
|
||||
portOfEndpoint(values, endpointPorts(m))
|
||||
composeName(values, r.PublicDomain, r.At, reaches, endpointPorts(m), blocks)
|
||||
out[to] = append(out[to], Contribution{From: m.Module, Values: values})
|
||||
}
|
||||
// Several contributions to one requirement (ADR 0094's sibling for `contributes`): an
|
||||
@@ -1112,21 +1135,11 @@ func (r Resolution) contributions(settings SettingsBy, grants []Grant,
|
||||
// name always reaches the provider from here.
|
||||
for _, to := range sortedKeys(m.ContributesMany) {
|
||||
for _, local := range sortedKeys(m.ContributesMany[to]) {
|
||||
values, err := settle(m.ContributesMany[to][local], settings[m.Module], nil,
|
||||
values, err := r.composed(m, m.ContributesMany[to][local], settings[m.Module],
|
||||
m.Module+" contributing "+local+" to "+to)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("%s contributing %s to %s: %w", m.Module, local, to, err)
|
||||
return nil, err
|
||||
}
|
||||
reaches, err := Reaches(m, settings[m.Module])
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("%s contributing %s to %s: %w", m.Module, local, to, err)
|
||||
}
|
||||
blocks, err := Endpoints(m, settings[m.Module])
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("%s contributing %s to %s: %w", m.Module, local, to, err)
|
||||
}
|
||||
portOfEndpoint(values, endpointPorts(m))
|
||||
composeName(values, r.PublicDomain, r.At, reaches, endpointPorts(m), blocks)
|
||||
out[to] = append(out[to], Contribution{From: m.Module, Values: values})
|
||||
}
|
||||
}
|
||||
@@ -1134,6 +1147,82 @@ func (r Resolution) contributions(settings SettingsBy, grants []Grant,
|
||||
return out, nil
|
||||
}
|
||||
|
||||
// composed is one contribution as its provider receives it: settled with this node's settings, its
|
||||
// endpoint's port filled in, and its names composed from the label.
|
||||
//
|
||||
// **One function, because two readers must agree.** The provider is told the names in its received
|
||||
// file; the contributing module is told the same names in its own binding (novox/hq 04-ISSUES/122).
|
||||
// Composing them twice, in two places, is how the proxy would come to serve one name while the
|
||||
// module wrote another into its configuration.
|
||||
func (r Resolution) composed(m Manifest, raw map[string]any, layers []Layer, what string) (
|
||||
map[string]any, error) {
|
||||
values, err := settle(raw, layers, nil, what)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("%s: %w", what, err)
|
||||
}
|
||||
reaches, err := Reaches(m, layers)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("%s: %w", what, err)
|
||||
}
|
||||
blocks, err := Endpoints(m, layers)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("%s: %w", what, err)
|
||||
}
|
||||
portOfEndpoint(values, endpointPorts(m))
|
||||
composeName(values, r.PublicDomain, r.At, reaches, endpointPorts(m), blocks)
|
||||
return values, nil
|
||||
}
|
||||
|
||||
// ownNames is what a module is known by through what it contributes to one requirement — the names
|
||||
// the mesh composed for it, and nothing else of the contribution.
|
||||
//
|
||||
// **The half a module could not learn** (novox/hq 04-ISSUES/122). A module contributes a label, the
|
||||
// mesh joins it with this node's domains, and the provider serves the result — and the module itself
|
||||
// was never told. Software that must know its own address (a login redirect, a canonical URL, an
|
||||
// issuer) had it written into the manifest as a literal, which is a domain in a definition and wrong
|
||||
// on every other machine. `${bound:<requirement>:name}` is the answer, from the same composition the
|
||||
// provider receives.
|
||||
//
|
||||
// Several contributions to one requirement are keyed by their local name under `names`.
|
||||
func (r Resolution) ownNames(m Manifest, to string, layers []Layer) (map[string]any, error) {
|
||||
pick := func(values map[string]any) map[string]any {
|
||||
names := map[string]any{}
|
||||
for _, key := range []string{"name", "internal-name"} {
|
||||
if v, ok := values[key].(string); ok && v != "" {
|
||||
names[key] = v
|
||||
}
|
||||
}
|
||||
return names
|
||||
}
|
||||
out := map[string]any{}
|
||||
if raw, ok := m.Contributes[to]; ok {
|
||||
values, err := r.composed(m, raw, layers, m.Module+" contributing to "+to)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
for k, v := range pick(values) {
|
||||
out[k] = v
|
||||
}
|
||||
}
|
||||
if locals := m.ContributesMany[to]; len(locals) > 0 {
|
||||
many := map[string]any{}
|
||||
for _, local := range sortedKeys(locals) {
|
||||
values, err := r.composed(m, locals[local], layers,
|
||||
m.Module+" contributing "+local+" to "+to)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if names := pick(values); len(names) > 0 {
|
||||
many[local] = names
|
||||
}
|
||||
}
|
||||
if len(many) > 0 {
|
||||
out["names"] = many
|
||||
}
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
|
||||
// composeName joins a contribution's label with a node's public domain, and separately with its
|
||||
// private one, in place (novox/hq ADR 0056).
|
||||
//
|
||||
@@ -1330,14 +1419,14 @@ func sortedKeys[V any](m map[string]V) []string {
|
||||
// Where it is and what the providing module said about using it. **No credential**, and the file
|
||||
// says so rather than leaving a reader to wonder whether one was meant to be there — a missing
|
||||
// field looks like a bug, and a stated absence looks like a boundary.
|
||||
func boundFile(n Needed, path, as string) (map[string]any, error) {
|
||||
func boundFile(n Needed, path, as string, own map[string]any) (map[string]any, error) {
|
||||
// A record has no machine and no address. Saying so is the difference between a reader
|
||||
// concluding "somewhere with no address" and concluding the mesh failed to fill something in.
|
||||
where := any(n.At)
|
||||
if n.ByRecord {
|
||||
where = "a record in this mesh, not a machine"
|
||||
}
|
||||
body, err := json.MarshalIndent(map[string]any{
|
||||
doc := map[string]any{
|
||||
"binding": 1,
|
||||
"provision": n.Name,
|
||||
"from": n.From,
|
||||
@@ -1353,7 +1442,14 @@ func boundFile(n Needed, path, as string) (map[string]any, error) {
|
||||
"generated": "by the mesh — do not edit; replaced whenever this changes. " +
|
||||
"The credential is not here: it is sealed, in the file this module's manifest " +
|
||||
"names under `secrets`",
|
||||
}, "", " ")
|
||||
}
|
||||
// **What this module is called through what it contributes here** (novox/hq 04-ISSUES/122):
|
||||
// `name`, `internal-name`, or `names` by local name — composed exactly as the provider receives
|
||||
// them. Absent when the module contributes nothing named, rather than written empty.
|
||||
for key, value := range own {
|
||||
doc[key] = value
|
||||
}
|
||||
body, err := json.MarshalIndent(doc, "", " ")
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
@@ -0,0 +1,130 @@
|
||||
package catalogue
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"strings"
|
||||
"testing"
|
||||
)
|
||||
|
||||
// A module that must know its own address — a login redirect, a canonical URL, an issuer — had it
|
||||
// written into its manifest as a literal (novox/hq 04-ISSUES/122): a domain in a definition, wrong on
|
||||
// every other machine. It is told instead, from the same composition the provider receives.
|
||||
|
||||
// selfAware contributes a labelled route, binds the requirement, and writes its own name into a file.
|
||||
func selfAware(label string) Manifest {
|
||||
m := labelled("board", label, 8080)
|
||||
m.Requires = []string{"reverse-proxy"}
|
||||
m.Binds = map[string]string{"reverse-proxy": "/var/lib/board/route.json"}
|
||||
m.Resources = []map[string]any{
|
||||
{"id": "conf", "type": "file", "path": "/var/lib/board/app.conf",
|
||||
"content": "root = https://${bound:reverse-proxy:name}/\ninternal = ${bound:reverse-proxy:internal-name}\n"},
|
||||
}
|
||||
return m
|
||||
}
|
||||
|
||||
// servingProxy is proxy() as the catalogue's route providers are declared: the provision scoped to
|
||||
// the mesh, serving nothing a consumer must know (route-adapter, route-proxy: `"serves": {"route": {}}`).
|
||||
func servingProxy() Manifest {
|
||||
p := proxy()
|
||||
p.Provides = []Offer{{Name: "reverse-proxy", Scope: ScopeMesh}}
|
||||
p.Serves = map[string]map[string]any{"reverse-proxy": {}}
|
||||
return p
|
||||
}
|
||||
|
||||
// nodeProxy is the same provider scoped to its node, whose answer on the same machine comes from
|
||||
// `here` rather than from the mesh's needs — the other path a binding is written by.
|
||||
func nodeProxy() Manifest {
|
||||
p := proxy()
|
||||
p.Serves = map[string]map[string]any{"reverse-proxy": {"scheme": "http"}}
|
||||
return p
|
||||
}
|
||||
|
||||
// onBoth is a node with a public domain and a private-network address, so both names compose.
|
||||
func onBoth(domain string) Node {
|
||||
n := withDomain(domain)
|
||||
n.At = "anchor.internal"
|
||||
return n
|
||||
}
|
||||
|
||||
func fileAt(t *testing.T, out []map[string]any, path string) string {
|
||||
t.Helper()
|
||||
for _, r := range out {
|
||||
if r["path"] == path {
|
||||
return r["content"].(string)
|
||||
}
|
||||
}
|
||||
t.Fatalf("nothing was declared at %s", path)
|
||||
return ""
|
||||
}
|
||||
|
||||
func TestAModuleIsToldTheNameItsProviderServes(t *testing.T) {
|
||||
got, err := Resolve(shelf(servingProxy(), selfAware("git")), []string{"traefik", "board"},
|
||||
onBoth("example.tld"), World{})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
out := mustDeclare(t, got)
|
||||
served := received(t, out)[0].Values
|
||||
|
||||
var binding map[string]any
|
||||
if err := json.Unmarshal([]byte(fileAt(t, out, "/var/lib/board/route.json")), &binding); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if binding["name"] != served["name"] || binding["name"] != "git.example.tld" {
|
||||
t.Fatalf("the module was told %v, the provider serves %v", binding["name"], served["name"])
|
||||
}
|
||||
if binding["internal-name"] != served["internal-name"] || binding["internal-name"] == nil {
|
||||
t.Fatalf("internal name: module told %v, provider serves %v",
|
||||
binding["internal-name"], served["internal-name"])
|
||||
}
|
||||
|
||||
conf := fileAt(t, out, "/var/lib/board/app.conf")
|
||||
want := "root = https://git.example.tld/\ninternal = " + served["internal-name"].(string) + "\n"
|
||||
if conf != want {
|
||||
t.Fatalf("the file was rendered as\n%s\nwant\n%s", conf, want)
|
||||
}
|
||||
}
|
||||
|
||||
func TestTheNameAModuleIsToldFollowsTheNodesDomain(t *testing.T) {
|
||||
// The whole point: the same definition, two machines, two names — nothing edited.
|
||||
for _, domain := range []string{"example.tld", "other.example"} {
|
||||
got, err := Resolve(shelf(servingProxy(), selfAware("git")), []string{"traefik", "board"},
|
||||
onBoth(domain), World{})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
conf := fileAt(t, mustDeclare(t, got), "/var/lib/board/app.conf")
|
||||
if !strings.HasPrefix(conf, "root = https://git."+domain+"/") {
|
||||
t.Fatalf("on %s the module wrote %q", domain, conf)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestAModuleWithNoPublicNameIsNotToldOne(t *testing.T) {
|
||||
// No public domain on the node: nothing composed, so no `name` — and a file asking for one is
|
||||
// refused rather than rendered with a placeholder or an empty host.
|
||||
m := selfAware("git")
|
||||
m.Resources[0]["content"] = "root = https://${bound:reverse-proxy:name}/\n"
|
||||
got, err := Resolve(shelf(servingProxy(), m), []string{"traefik", "board"}, workstation(), World{})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := got.Declaration(Rendering{}); err == nil ||
|
||||
!strings.Contains(err.Error(), `"name"`) {
|
||||
t.Fatalf("a file asking for a name that was never composed was not refused: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestAModuleIsToldItsNameByANodeScopedProviderToo(t *testing.T) {
|
||||
m := selfAware("git")
|
||||
m.Resources[0]["content"] = "root = https://${bound:reverse-proxy:name}/\n"
|
||||
got, err := Resolve(shelf(nodeProxy(), m), []string{"board"},
|
||||
withDomain("example.tld"), World{})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
conf := fileAt(t, mustDeclare(t, got), "/var/lib/board/app.conf")
|
||||
if !strings.HasPrefix(conf, "root = https://git.example.tld/") {
|
||||
t.Fatalf("a same-machine, node-scoped answer did not tell the module its name: %q", conf)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user