Compare commits
1
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
24f024dd74 |
@@ -404,7 +404,8 @@ func declarationWith(ctx context.Context, open *stores, node string,
|
||||
if err != nil {
|
||||
return sendable{}, err
|
||||
}
|
||||
return sendable{Resources: composed.Resources, Adoption: adoption}, nil
|
||||
return sendable{Resources: composed.Resources, Adoption: adoption,
|
||||
Received: composed.Received, Mesh: with.Mesh}, nil
|
||||
}
|
||||
|
||||
// renderingFor is everything a node's declaration is composed with, and the node's record.
|
||||
|
||||
+23
-10
@@ -389,11 +389,7 @@ func pushCommand(ctx context.Context, args []string) error {
|
||||
fmt.Printf("\n%d node(s) told\n", len(sending))
|
||||
// And each machine's memberships, as every other send does (ADR 0160): a push is the one most
|
||||
// operators run, and on 2026-10-01 it was the one path that issued none.
|
||||
var told []string
|
||||
for _, s := range sending {
|
||||
told = append(told, s.node)
|
||||
}
|
||||
if err := issueMemberships(ctx, open, server, told); err != nil {
|
||||
if err := issueMemberships(ctx, open, server, sending); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
@@ -706,11 +702,15 @@ func sendTo(ctx context.Context, open *stores, names []string) error {
|
||||
// And every assignment on those machines its membership (novox/hq ADR 0160): composed from the
|
||||
// same records the bus's accounts are, so what a runtime serves and what its account may are one
|
||||
// composition. Issued after the declaration, because the runtime it is for arrives with it.
|
||||
return issueMemberships(ctx, open, server, names)
|
||||
return issueMemberships(ctx, open, server, sending)
|
||||
}
|
||||
|
||||
// issueMemberships publishes the membership of every module on the named machines.
|
||||
func issueMemberships(ctx context.Context, open *stores, server *link.Server, names []string) error {
|
||||
// issueMemberships publishes the membership of every module on the machines just sent.
|
||||
//
|
||||
// Each carries what its module receives and the private network's addresses, from the same
|
||||
// composition as the declaration it was sent (novox/hq ADR 0167): a provider reads what it is
|
||||
// given on the bus, and the file written beside it says the same thing.
|
||||
func issueMemberships(ctx context.Context, open *stores, server *link.Server, sent []readyNode) error {
|
||||
records, err := open.inventory.BusRecords(ctx)
|
||||
if err != nil {
|
||||
return err
|
||||
@@ -725,9 +725,22 @@ func issueMemberships(ctx context.Context, open *stores, server *link.Server, na
|
||||
// the push stands, the first failure is named once, and the next push tries again.
|
||||
issued, failed := 0, 0
|
||||
var first error
|
||||
for _, node := range names {
|
||||
for _, s := range sent {
|
||||
node := s.node
|
||||
for _, d := range records.Assigned[node] {
|
||||
body, err := json.Marshal(broker.MembershipFor(node, d, where))
|
||||
membership := broker.MembershipFor(node, d, where)
|
||||
membership.Mesh = s.declared.Mesh
|
||||
for requirement, given := range s.declared.Received[d.Module] {
|
||||
raw, err := json.Marshal(given)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if membership.Receives == nil {
|
||||
membership.Receives = map[string]json.RawMessage{}
|
||||
}
|
||||
membership.Receives[requirement] = raw
|
||||
}
|
||||
body, err := json.Marshal(membership)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
@@ -25,6 +25,12 @@ type sendable struct {
|
||||
// Adoption is nil for a converged node, and then the body is byte for byte what it was before
|
||||
// adoption existed: an older host parses the envelope strictly and would refuse the key.
|
||||
Adoption *adoptionEnvelope
|
||||
|
||||
// Received and Mesh are not sent in the declaration. They are what this machine's memberships
|
||||
// are issued with on the bus (novox/hq ADR 0167): each module's received contributions, from
|
||||
// the same composition as its received files, and every machine's private-network address.
|
||||
Received map[string]map[string][]catalogue.Contribution
|
||||
Mesh []string
|
||||
}
|
||||
|
||||
// adoptionEnvelope is what an adopted node is told about its mode. Taken is every module taken on
|
||||
|
||||
@@ -0,0 +1,132 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"log"
|
||||
"os"
|
||||
"strings"
|
||||
"sync/atomic"
|
||||
"time"
|
||||
|
||||
"github.com/nats-io/nats.go"
|
||||
|
||||
"github.com/novox/mesh-controller/internal/broker"
|
||||
)
|
||||
|
||||
// What the mesh issued this proxy, read on the bus (novox/hq ADR 0160, ADR 0167).
|
||||
//
|
||||
// **The proxy is told, not left to work it out.** Its membership carries the routes it is given —
|
||||
// the same contributions its file is written from — and every machine's address on the private
|
||||
// network, which is who may be served an internal name. Read once at connect and followed, so a
|
||||
// route added or a machine joining reaches a running proxy without a restart.
|
||||
|
||||
// credential is the bus account the mesh delivered as this module's own secret named broker.
|
||||
type credential struct {
|
||||
URL string `json:"url"`
|
||||
Fingerprint string `json:"fingerprint"`
|
||||
Node string `json:"node"`
|
||||
Module string `json:"module"`
|
||||
User string `json:"user"`
|
||||
Password string `json:"password"`
|
||||
}
|
||||
|
||||
// followMembership connects with the credential in path and applies every membership the mesh
|
||||
// issues this proxy. It retries the first connection for as long as it takes: a proxy that started
|
||||
// before the bus keeps serving the file, and takes the bus when it answers.
|
||||
func followMembership(path string, held *table, fromBus *atomic.Bool) {
|
||||
for {
|
||||
err := followOnce(path, held, fromBus)
|
||||
if err == nil {
|
||||
return
|
||||
}
|
||||
log.Printf("cannot follow this proxy's membership, serving the file meanwhile: %v", err)
|
||||
time.Sleep(30 * time.Second)
|
||||
}
|
||||
}
|
||||
|
||||
func followOnce(path string, held *table, fromBus *atomic.Bool) error {
|
||||
raw, err := os.ReadFile(path)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
var cred credential
|
||||
if err := json.Unmarshal(raw, &cred); err != nil {
|
||||
return fmt.Errorf("the broker credential is not one: %w", err)
|
||||
}
|
||||
if cred.Node == "" || cred.Module == "" {
|
||||
return fmt.Errorf("the broker credential names no node or module, so it has no membership")
|
||||
}
|
||||
|
||||
opts := []nats.Option{
|
||||
nats.Name(cred.Node + "." + cred.Module),
|
||||
nats.UserInfo(cred.User, cred.Password),
|
||||
// Its own inbox, and nothing wider: every principal is granted `_INBOX.<its user>.>` alone.
|
||||
nats.CustomInboxPrefix("_INBOX." + cred.User),
|
||||
// The bus being restarted is an upgrade, not a reason to stop following.
|
||||
nats.MaxReconnects(-1),
|
||||
}
|
||||
if strings.TrimSpace(cred.Fingerprint) != "" {
|
||||
opts = append(opts, nats.Secure(broker.PinnedToFingerprint(cred.Fingerprint)))
|
||||
}
|
||||
conn, err := nats.Connect(cred.URL, opts...)
|
||||
if err != nil {
|
||||
return fmt.Errorf("connecting to the bus at %s: %w", broker.BareAddress(cred.URL), err)
|
||||
}
|
||||
|
||||
subject := broker.MembershipSubject(cred.Node, cred.Module)
|
||||
apply := func(body []byte) {
|
||||
var issued broker.Membership
|
||||
if err := json.Unmarshal(body, &issued); err != nil {
|
||||
log.Printf("a membership arrived that is not one: %v", err)
|
||||
return
|
||||
}
|
||||
if took := applyMembership(issued, held); took && !fromBus.Swap(true) {
|
||||
log.Printf("routes now come from this proxy's membership on %s", subject)
|
||||
}
|
||||
}
|
||||
|
||||
// Followed first, read second: an issue landing between the two is applied, not missed.
|
||||
if _, err := conn.Subscribe(subject, func(m *nats.Msg) { apply(m.Data) }); err != nil {
|
||||
conn.Close()
|
||||
return fmt.Errorf("cannot follow %s: %w", subject, err)
|
||||
}
|
||||
// The subject-addressed direct get: the one request this account may make of the stream.
|
||||
got, err := conn.Request("$JS.API.DIRECT.GET."+broker.AssignmentsStream+"."+subject, nil, 5*time.Second)
|
||||
switch {
|
||||
case err != nil:
|
||||
log.Printf("cannot read the membership issued on %s yet (%v); following it", subject, err)
|
||||
case got.Header.Get("Status") != "" || len(got.Data) == 0:
|
||||
log.Printf("no membership issued on %s yet; serving the file until one is", subject)
|
||||
default:
|
||||
apply(got.Data)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// applyMembership serves what a membership says, and says whether it said anything about routes.
|
||||
//
|
||||
// A membership with no routes in it is one from a controller older than ADR 0167, and the file stays
|
||||
// the source rather than every route being withdrawn because a field was absent.
|
||||
func applyMembership(issued broker.Membership, held *table) bool {
|
||||
raw, carries := issued.Receives["route"]
|
||||
if !carries {
|
||||
return false
|
||||
}
|
||||
var contributions []contribution
|
||||
if err := json.Unmarshal(raw, &contributions); err != nil {
|
||||
log.Printf("the routes in this proxy's membership are not contributions, keeping what is served: %v", err)
|
||||
return false
|
||||
}
|
||||
inside, err := sourcesOf(issued.Mesh)
|
||||
if err != nil {
|
||||
log.Printf("the mesh in this proxy's membership is unreadable, keeping what is served: %v", err)
|
||||
return false
|
||||
}
|
||||
routes, public := routesOf(contributions)
|
||||
held.set(routes, public)
|
||||
held.setInside(inside)
|
||||
log.Printf("serving %d route(s) from the membership, internal names to %d machine(s): %s",
|
||||
len(routes), len(inside), strings.Join(held.names(), ", "))
|
||||
return true
|
||||
}
|
||||
@@ -61,6 +61,7 @@ import (
|
||||
"sort"
|
||||
"strings"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"time"
|
||||
|
||||
"golang.org/x/crypto/acme"
|
||||
@@ -197,77 +198,44 @@ type table struct {
|
||||
// pass ACME's own validation (it has no public DNS to prove it against), so asking for it is
|
||||
// not merely pointless but the failing order onlyWhatTheMeshSaid exists to prevent.
|
||||
public map[string]bool
|
||||
// inside is the private network's range, where a request must come from to be served a name
|
||||
// that is only internal. Set once at start, never replaced with the routes: it is what the
|
||||
// private network is, not what is routed on it.
|
||||
// inside is where a request must come from to be served a name that is only internal: every
|
||||
// machine's address on the private network, as the mesh issued it in this proxy's membership
|
||||
// (novox/hq ADR 0167). Empty until it is issued, and then only the machine itself is inside.
|
||||
inside sources
|
||||
// bridges is the machine's own container networks, read from its interfaces and refreshed with
|
||||
// the routes, since a compose network can appear at any time.
|
||||
bridges sources
|
||||
}
|
||||
|
||||
// sources is the private network, as address ranges. The machine itself is always inside it —
|
||||
// anything on a machine may call anything on it (novox/hq ADR 0144) — so loopback needs no range.
|
||||
// sources is who may be served an internal name: the private network's addresses as the mesh
|
||||
// issued them. The machine itself is always inside — anything on a machine may call anything on it
|
||||
// (novox/hq ADR 0144) — so loopback needs no entry.
|
||||
type sources []netip.Prefix
|
||||
|
||||
// sourcesFrom reads the ranges the mesh wrote, separated by commas or spaces. A range that does not
|
||||
// parse is an error, not a range skipped: the proxy would otherwise serve internal names to fewer
|
||||
// machines than the mesh said, or start believing a typo.
|
||||
func sourcesFrom(text string) (sources, error) {
|
||||
// sourcesOf reads the addresses the mesh issued, each a single address or a range. One that does
|
||||
// not parse is an error, not an entry skipped: the proxy would otherwise serve internal names to
|
||||
// fewer machines than the mesh said, and say nothing.
|
||||
func sourcesOf(mesh []string) (sources, error) {
|
||||
var out sources
|
||||
for _, field := range strings.FieldsFunc(text, func(r rune) bool { return r == ',' || r == ' ' || r == '\n' || r == '\t' }) {
|
||||
prefix, err := netip.ParsePrefix(field)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("%q is not an address range: %w", field, err)
|
||||
for _, entry := range mesh {
|
||||
entry = strings.TrimSpace(entry)
|
||||
if prefix, err := netip.ParsePrefix(entry); err == nil {
|
||||
out = append(out, prefix.Masked())
|
||||
continue
|
||||
}
|
||||
out = append(out, prefix.Masked())
|
||||
addr, err := netip.ParseAddr(entry)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("%q is not an address on the private network", entry)
|
||||
}
|
||||
addr = addr.Unmap()
|
||||
out = append(out, netip.PrefixFrom(addr, addr.BitLen()))
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
|
||||
// bridgesFrom is the address ranges of the machine's container bridges — the same interfaces the
|
||||
// mesh's guard names as the machine itself (docker0, and the br-* a compose network gets), so the
|
||||
// proxy and the guard agree on what "this machine" is (novox/hq ADR 0144).
|
||||
func bridgesFrom(interfaces map[string][]net.Addr) sources {
|
||||
var out sources
|
||||
for name, addrs := range interfaces {
|
||||
if name != "docker0" && !strings.HasPrefix(name, "br-") {
|
||||
continue
|
||||
}
|
||||
for _, a := range addrs {
|
||||
if ipnet, ok := a.(*net.IPNet); ok {
|
||||
if prefix, err := netip.ParsePrefix(ipnet.String()); err == nil {
|
||||
out = append(out, prefix.Masked())
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
// theseBridges reads this machine's interfaces for bridgesFrom. An interface that cannot be read
|
||||
// contributes nothing: fewer callers inside, never more.
|
||||
func theseBridges() sources {
|
||||
interfaces, err := net.Interfaces()
|
||||
if err != nil {
|
||||
return nil
|
||||
}
|
||||
named := map[string][]net.Addr{}
|
||||
for _, i := range interfaces {
|
||||
if addrs, err := i.Addrs(); err == nil {
|
||||
named[i.Name] = addrs
|
||||
}
|
||||
}
|
||||
return bridgesFrom(named)
|
||||
}
|
||||
|
||||
// holds says whether a request from this remote address came from inside these ranges, or from
|
||||
// the machine itself.
|
||||
// holds says whether a request from this remote address came from the mesh or the machine itself.
|
||||
//
|
||||
// **By source, which the guard deliberately is not** — it names interfaces because a source
|
||||
// **By source, which the mesh's guard deliberately is not** — it names interfaces, because a source
|
||||
// address can be claimed by whoever sends the packet. The proxy cannot see the interface a request
|
||||
// arrived on, and here the claim does not carry: a connection needs its replies, and replies to a
|
||||
// mesh or container address leave by the tunnel or a local bridge, never back to the claimant.
|
||||
// mesh address leave by the tunnel, never back to the claimant.
|
||||
func (s sources) holds(remote string) bool {
|
||||
host := remote
|
||||
if h, _, err := net.SplitHostPort(remote); err == nil {
|
||||
@@ -418,13 +386,13 @@ func (t *table) hiddenFrom(host, remote string) bool {
|
||||
}
|
||||
t.mu.RLock()
|
||||
defer t.mu.RUnlock()
|
||||
return !t.inside.holds(remote) && !t.bridges.holds(remote)
|
||||
return !t.inside.holds(remote)
|
||||
}
|
||||
|
||||
// setBridges replaces the machine's container networks.
|
||||
func (t *table) setBridges(bridges sources) {
|
||||
// setInside replaces who the mesh is, as the membership said.
|
||||
func (t *table) setInside(inside sources) {
|
||||
t.mu.Lock()
|
||||
t.bridges = bridges
|
||||
t.inside = inside
|
||||
t.mu.Unlock()
|
||||
}
|
||||
|
||||
@@ -469,19 +437,20 @@ func run() error {
|
||||
}
|
||||
|
||||
held := newTable()
|
||||
// Unset means only this machine and its containers are inside, which serves an internal-only
|
||||
// name to nobody else — refused rather than served to everyone, which is what the proxy did
|
||||
// before it knew.
|
||||
inside, err := sourcesFrom(os.Getenv("INTERNAL_SOURCES"))
|
||||
if err != nil {
|
||||
return fmt.Errorf("INTERNAL_SOURCES: %w", err)
|
||||
// **The bus first, the file until it has spoken** (novox/hq ADR 0167). The membership carries
|
||||
// the routes and who the mesh is; the file carries the routes alone, so while the proxy reads
|
||||
// it an internal name is served to this machine and to nobody else — refused, never opened.
|
||||
fromBus := &atomic.Bool{}
|
||||
if credential := strings.TrimSpace(os.Getenv("MESH_BROKER_FILE")); credential != "" {
|
||||
go followMembership(credential, held, fromBus)
|
||||
} else {
|
||||
log.Printf("MESH_BROKER_FILE is not set: routes come from %s alone, and a name that is only "+
|
||||
"internal is served to this machine alone", path)
|
||||
}
|
||||
if len(inside) == 0 {
|
||||
log.Printf("INTERNAL_SOURCES is not set: a name that is only internal is served to this machine " +
|
||||
"and its containers alone")
|
||||
}
|
||||
held.inside = inside
|
||||
read := func() {
|
||||
if fromBus.Load() {
|
||||
return
|
||||
}
|
||||
routes, public, err := routesFrom(path)
|
||||
if err != nil {
|
||||
// Kept serving what it had. A file being rewritten is momentarily unreadable, and
|
||||
@@ -491,7 +460,6 @@ func run() error {
|
||||
return
|
||||
}
|
||||
held.set(routes, public)
|
||||
held.setBridges(theseBridges())
|
||||
log.Printf("serving %d route(s): %s", len(routes), strings.Join(held.names(), ", "))
|
||||
}
|
||||
read()
|
||||
@@ -886,10 +854,16 @@ func routesFrom(path string) (map[string][]rule, map[string]bool, error) {
|
||||
if err := json.Unmarshal(raw, &said); err != nil {
|
||||
return nil, nil, err
|
||||
}
|
||||
routes, public := routesOf(said.Given)
|
||||
return routes, public, nil
|
||||
}
|
||||
|
||||
// routesOf turns what the mesh gave into host → the rules for that host, and which hosts are public
|
||||
// names — the same whether the contributions came in the file or in the membership.
|
||||
func routesOf(contributions []contribution) (map[string][]rule, map[string]bool) {
|
||||
out := map[string][]rule{}
|
||||
public := map[string]bool{}
|
||||
for _, c := range said.Given {
|
||||
for _, c := range contributions {
|
||||
name, _ := c.Values["name"].(string)
|
||||
name = strings.TrimSpace(name)
|
||||
internal, _ := c.Values["internal-name"].(string)
|
||||
@@ -994,7 +968,7 @@ func routesFrom(path string) (map[string][]rule, map[string]bool, error) {
|
||||
out[strings.ToLower(internal)] = append(out[strings.ToLower(internal)], made)
|
||||
}
|
||||
}
|
||||
return out, public, nil
|
||||
return out, public
|
||||
}
|
||||
|
||||
// asWhole is any whole number the mesh wrote, whatever its magnitude.
|
||||
|
||||
@@ -2,6 +2,7 @@ package main
|
||||
|
||||
import (
|
||||
"crypto/tls"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"io"
|
||||
"net"
|
||||
@@ -10,11 +11,13 @@ import (
|
||||
"net/url"
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
"github.com/novox/mesh-controller/internal/broker"
|
||||
)
|
||||
|
||||
// behind is a workload the proxy can send to, and a table routing one public name and one
|
||||
// internal-only name to it, with the private network set to inside.
|
||||
func behind(t *testing.T, inside string) *table {
|
||||
// internal-only name to it, with the mesh's machines as the membership would issue them.
|
||||
func behind(t *testing.T, mesh ...string) *table {
|
||||
t.Helper()
|
||||
workload := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
io.WriteString(w, "the workload")
|
||||
@@ -33,10 +36,11 @@ func behind(t *testing.T, inside string) *table {
|
||||
t.Fatal(err)
|
||||
}
|
||||
held := newTable()
|
||||
held.inside, err = sourcesFrom(inside)
|
||||
inside, err := sourcesOf(mesh)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
held.setInside(inside)
|
||||
held.set(routes, public)
|
||||
return held
|
||||
}
|
||||
@@ -54,7 +58,7 @@ func askFrom(held *table, host, remote string) (int, string) {
|
||||
// 0138, issue 191). The proxy answers public names on the same listeners, so without this a name
|
||||
// being internal kept nobody out: a request from the internet only had to carry it.
|
||||
func TestAnInternalOnlyNameIsServedOnlyInsideThePrivateNetwork(t *testing.T) {
|
||||
held := behind(t, "10.10.0.0/24")
|
||||
held := behind(t, "10.10.0.1", "10.10.0.7")
|
||||
|
||||
if code, body := askFrom(held, "admin.anchor.internal", "10.10.0.7:51000"); code != http.StatusOK ||
|
||||
body != "the workload" {
|
||||
@@ -83,7 +87,7 @@ func TestAnInternalOnlyNameIsServedOnlyInsideThePrivateNetwork(t *testing.T) {
|
||||
// an outsider only under the public name. Nothing is lost — the outsider has the public name — and a
|
||||
// name stays one thing whichever route it came from.
|
||||
func TestAnInternalAliasOfAPublicRouteIsServedInsideOnly(t *testing.T) {
|
||||
held := behind(t, "10.10.0.0/24")
|
||||
held := behind(t, "10.10.0.1", "10.10.0.7")
|
||||
if code, body := askFrom(held, "app.anchor.internal", "10.10.0.7:51000"); code != http.StatusOK {
|
||||
t.Errorf("the internal alias stopped answering the private network: %d %q", code, body)
|
||||
}
|
||||
@@ -95,32 +99,10 @@ func TestAnInternalAliasOfAPublicRouteIsServedInsideOnly(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
// A container on this machine reaches the proxy from its bridge's range, and is the machine itself
|
||||
// (novox/hq ADR 0144): inside, though it is neither loopback nor the mesh.
|
||||
func TestAContainerOnThisMachineIsInside(t *testing.T) {
|
||||
held := behind(t, "10.10.0.0/24")
|
||||
_, bridge, _ := net.ParseCIDR("172.18.0.1/16")
|
||||
bridge.IP = net.ParseIP("172.18.0.1")
|
||||
_, other, _ := net.ParseCIDR("192.168.1.20/24")
|
||||
other.IP = net.ParseIP("192.168.1.20")
|
||||
held.setBridges(bridgesFrom(map[string][]net.Addr{
|
||||
"br-0123456789ab": {bridge},
|
||||
"eth0": {other},
|
||||
}))
|
||||
|
||||
if code, _ := askFrom(held, "admin.anchor.internal", "172.18.0.5:51000"); code != http.StatusOK {
|
||||
t.Errorf("a container on this machine was refused: %d", code)
|
||||
}
|
||||
// The machine's own network is not a container bridge: a neighbour there is not the machine.
|
||||
if code, _ := askFrom(held, "admin.anchor.internal", "192.168.1.30:51000"); code != http.StatusNotFound {
|
||||
t.Errorf("a neighbour on the machine's network was served an internal-only name: %d", code)
|
||||
}
|
||||
}
|
||||
|
||||
// With no private network said, only the machine itself is inside — refused to everyone else,
|
||||
// never served to everyone.
|
||||
func TestWithNoPrivateNetworkSaidAnInternalOnlyNameIsServedToTheMachineAlone(t *testing.T) {
|
||||
held := behind(t, "")
|
||||
// Before a membership has said who the mesh is, only the machine itself is inside — refused to
|
||||
// everyone else, never served to everyone.
|
||||
func TestUntilTheMeshIsIssuedAnInternalOnlyNameIsServedToTheMachineAlone(t *testing.T) {
|
||||
held := behind(t)
|
||||
if code, _ := askFrom(held, "admin.anchor.internal", "10.10.0.7:51000"); code != http.StatusNotFound {
|
||||
t.Errorf("an internal-only name was served with no private network said: %d", code)
|
||||
}
|
||||
@@ -139,7 +121,7 @@ func (c from) RemoteAddr() net.Addr { return c.remote }
|
||||
// The handshake refuses an internal-only name to an outsider too: the certificate would name it,
|
||||
// and serving it would answer the question the routing refuses to.
|
||||
func TestTheHandshakeRefusesAnInternalOnlyNameToAnOutsider(t *testing.T) {
|
||||
held := behind(t, "10.10.0.0/24")
|
||||
held := behind(t, "10.10.0.1", "10.10.0.7")
|
||||
served := &tls.Certificate{}
|
||||
pick := certificateFor(held, func(*tls.ClientHelloInfo) (*tls.Certificate, error) { return served, nil }, nil)
|
||||
hello := func(name, remote string) *tls.ClientHelloInfo {
|
||||
@@ -158,26 +140,67 @@ func TestTheHandshakeRefusesAnInternalOnlyNameToAnOutsider(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
// A range the proxy cannot read stops it, rather than serving internal names to fewer machines
|
||||
// than the mesh said, or to a typo.
|
||||
func TestAPrivateNetworkThatDoesNotParseIsRefused(t *testing.T) {
|
||||
if _, err := sourcesFrom("10.10.0.0/24, not-a-range"); err == nil {
|
||||
t.Error("a range that does not parse was accepted")
|
||||
// The mesh is issued as machines' addresses; a range is read as well. One that does not parse is
|
||||
// refused rather than skipped, so a typo never quietly narrows or widens who is inside.
|
||||
func TestTheMeshIsReadAsAddressesAndRanges(t *testing.T) {
|
||||
if _, err := sourcesOf([]string{"10.10.0.1", "not-an-address"}); err == nil {
|
||||
t.Error("an entry that is not an address was accepted")
|
||||
}
|
||||
inside, err := sourcesFrom("10.10.0.0/24 fd00::/8")
|
||||
inside, err := sourcesOf([]string{"10.10.0.1", "fd00::1", "10.20.0.0/24"})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
for remote, want := range map[string]bool{
|
||||
"10.10.0.200:1": true,
|
||||
"[::ffff:10.10.0.3]:1": true,
|
||||
"10.10.0.1:1": true,
|
||||
"[::ffff:10.10.0.1]:1": true,
|
||||
"[fd00::1]:1": true,
|
||||
"10.11.0.1:1": false,
|
||||
"10.20.0.200:1": true,
|
||||
"10.10.0.2:1": false,
|
||||
"192.168.1.10:1": false,
|
||||
"not-an-address": false,
|
||||
} {
|
||||
if inside.holds(remote) != want {
|
||||
t.Errorf("%s inside the private network: got %v, want %v", remote, !want, want)
|
||||
t.Errorf("%s inside the mesh: got %v, want %v", remote, !want, want)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// What the mesh issues is what is served: the routes in the membership, internal names to the
|
||||
// machines it names (novox/hq ADR 0167).
|
||||
func TestAMembershipIsServedAsIssued(t *testing.T) {
|
||||
held := newTable()
|
||||
took := applyMembership(broker.Membership{
|
||||
Receives: map[string]json.RawMessage{"route": json.RawMessage(`[
|
||||
{"from":"admin","node":"anchor","at":"anchor.internal",
|
||||
"values":{"internal-name":"admin.anchor.internal","port":8080}}]`)},
|
||||
Mesh: []string{"10.10.0.7"},
|
||||
}, held)
|
||||
if !took {
|
||||
t.Fatal("a membership carrying routes was not applied")
|
||||
}
|
||||
if code, _ := askFrom(held, "admin.anchor.internal", "10.10.0.7:1"); code == http.StatusNotFound {
|
||||
t.Error("a machine the membership names was refused the internal-only route")
|
||||
}
|
||||
if code, _ := askFrom(held, "admin.anchor.internal", "10.10.0.9:1"); code != http.StatusNotFound {
|
||||
t.Errorf("a machine the membership does not name was served the internal-only route: %d", code)
|
||||
}
|
||||
}
|
||||
|
||||
// A membership that says nothing about routes is one from a controller that does not issue them,
|
||||
// and changes nothing: the file stays the source rather than every route being withdrawn.
|
||||
func TestAMembershipWithoutRoutesLeavesTheFileServing(t *testing.T) {
|
||||
held := behind(t, "10.10.0.7")
|
||||
before := held.names()
|
||||
if applyMembership(broker.Membership{Mesh: []string{"10.10.0.7"}}, held) {
|
||||
t.Error("a membership without routes was taken as the source of routes")
|
||||
}
|
||||
if got := held.names(); strings.Join(got, ",") != strings.Join(before, ",") {
|
||||
t.Errorf("a membership without routes changed what is served: %v, was %v", got, before)
|
||||
}
|
||||
if applyMembership(broker.Membership{
|
||||
Receives: map[string]json.RawMessage{"route": json.RawMessage(`[]`)},
|
||||
Mesh: []string{"not-an-address"},
|
||||
}, held) {
|
||||
t.Error("a membership whose mesh cannot be read was applied")
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
package broker
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"sort"
|
||||
"strings"
|
||||
)
|
||||
@@ -33,6 +34,17 @@ type Membership struct {
|
||||
Reaches map[string][]string `json:"reaches,omitempty"`
|
||||
// Tools is where this instance answers what it serves — the runtime's one verb of its own.
|
||||
Tools string `json:"tools"`
|
||||
// Receives is what this assignment is given for each requirement it receives, by requirement:
|
||||
// the contributions of every module that asked for it, as the catalogue composed them (novox/hq
|
||||
// ADR 0167). The same list its received file is written from, so the two cannot disagree; a
|
||||
// requirement nobody contributed to is an empty list, never absent. Kept as JSON because the
|
||||
// catalogue owns the shape of a contribution and the bus only carries it.
|
||||
Receives map[string]json.RawMessage `json:"receives,omitempty"`
|
||||
// Mesh is every machine's address on the private network — what a rule saying "from the mesh"
|
||||
// resolves to in the packet filter, issued here from the same list (novox/hq ADR 0167). A
|
||||
// module that must tell the mesh from the world, the route proxy serving an internal name, reads
|
||||
// it here rather than keeping a definition of its own.
|
||||
Mesh []string `json:"mesh,omitempty"`
|
||||
}
|
||||
|
||||
// Served is one address a tool is answered on.
|
||||
|
||||
@@ -241,15 +241,21 @@ func (r Resolution) Declaration(with Rendering) ([]map[string]any, error) {
|
||||
// Owner is kept beside the resources because a resource id cannot be split back into its module:
|
||||
// a module's name may itself contain a dot. What the mesh adds of its own — an opening, the guard —
|
||||
// has no owner.
|
||||
//
|
||||
// Received is what each module on the machine is given for each requirement it receives — the same
|
||||
// contributions its received file is written from, kept beside it so the mesh can also issue them
|
||||
// on the bus in the module's membership (novox/hq ADR 0167). By module, then requirement.
|
||||
type Composed struct {
|
||||
Resources []map[string]any
|
||||
Owner map[string]string
|
||||
Received map[string]map[string][]Contribution
|
||||
}
|
||||
|
||||
// Compose is Declaration with the owner of every resource said.
|
||||
func (r Resolution) Compose(with Rendering) (Composed, error) {
|
||||
owner := map[string]string{}
|
||||
resources, err := r.compose(with, owner)
|
||||
received := map[string]map[string][]Contribution{}
|
||||
resources, err := r.compose(with, owner, received)
|
||||
if err != nil {
|
||||
return Composed{}, err
|
||||
}
|
||||
@@ -261,7 +267,7 @@ func (r Resolution) Compose(with Rendering) (Composed, error) {
|
||||
"sealed": with.BusMembership, "mode": "0600",
|
||||
})
|
||||
}
|
||||
return Composed{Resources: resources, Owner: owner}, nil
|
||||
return Composed{Resources: resources, Owner: owner, Received: received}, nil
|
||||
}
|
||||
|
||||
// BusMembershipID names the resource carrying a machine's membership for the new bus, and
|
||||
@@ -270,7 +276,8 @@ func BusMembershipID() string { return "bus-membership" }
|
||||
|
||||
const BusMembershipPath = "/var/lib/mesh/membership-next.json"
|
||||
|
||||
func (r Resolution) compose(with Rendering, owner map[string]string) ([]map[string]any, error) {
|
||||
func (r Resolution) compose(with Rendering, owner map[string]string,
|
||||
received map[string]map[string][]Contribution) ([]map[string]any, error) {
|
||||
// Every manifest is placed first (novox/hq ADR 0112): the maps naming where its bindings,
|
||||
// credentials and contributions land are resolved against this node's directories, so every
|
||||
// reader below — the binding files, the sealed secrets, the grant paths a contribution
|
||||
@@ -596,6 +603,12 @@ func (r Resolution) compose(with Rendering, owner map[string]string) ([]map[stri
|
||||
return nil, err
|
||||
}
|
||||
first = append(first, file)
|
||||
if received[m.Module] == nil {
|
||||
received[m.Module] = map[string][]Contribution{}
|
||||
}
|
||||
// Empty rather than absent when nobody contributed, for the reason the file is
|
||||
// written empty: "nothing asked" and "never told" want different responses.
|
||||
received[m.Module][to] = append([]Contribution{}, given[to]...)
|
||||
}
|
||||
if m.Keeps != "" && with.Kept != nil {
|
||||
file, err := keptFile(m.Keeps, with.Kept)
|
||||
|
||||
@@ -0,0 +1,66 @@
|
||||
package catalogue
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"reflect"
|
||||
"testing"
|
||||
)
|
||||
|
||||
// What a provider receives is composed once, and issued twice: as its received file, and in its
|
||||
// membership on the bus (novox/hq ADR 0167). The two are the same list, so a proxy reading the bus
|
||||
// and one reading the file serve the same routes — including the port the machine published, which
|
||||
// is the same-node fix the file already carries.
|
||||
func TestWhatAProviderReceivesIsTheSameOnTheBusAsInItsFile(t *testing.T) {
|
||||
gitea := Manifest{
|
||||
Module: "gitea", Version: "1",
|
||||
Listens: []Listening{{Port: 3000, Protocol: "tcp", From: FromMesh}},
|
||||
Contributes: map[string]map[string]any{"route": {"label": "git", "port": 3000}},
|
||||
Resources: []map[string]any{{
|
||||
"id": "server", "type": "container", "name": "gitea", "ports": []any{"3000"},
|
||||
}},
|
||||
}
|
||||
r, err := Resolve(shelf(gitea, routeProxy(), stepCA()),
|
||||
[]string{"gitea", "route-proxy", "step-ca"}, reachable(), World{})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
composed, err := r.Compose(Rendering{Ports: map[string]map[int]int{"gitea": {3000: 20000}}})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
file := fileNamed(composed.Resources, "route-proxy.received-route")
|
||||
if file == nil {
|
||||
t.Fatal("the proxy was given no routes file")
|
||||
}
|
||||
var written struct {
|
||||
Given []Contribution `json:"given"`
|
||||
}
|
||||
if err := json.Unmarshal([]byte(file["content"].(string)), &written); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
issued, said := composed.Received["route-proxy"]["route"]
|
||||
if !said {
|
||||
t.Fatalf("nothing is issued for the proxy to receive on the bus: %v", composed.Received)
|
||||
}
|
||||
// Compared as JSON, which is what both are once they leave the controller.
|
||||
a, _ := json.Marshal(written.Given)
|
||||
b, _ := json.Marshal(issued)
|
||||
var fromFile, fromBus any
|
||||
_ = json.Unmarshal(a, &fromFile)
|
||||
_ = json.Unmarshal(b, &fromBus)
|
||||
if !reflect.DeepEqual(fromFile, fromBus) {
|
||||
t.Errorf("the bus and the file disagree about the routes:\nfile %s\nbus %s", a, b)
|
||||
}
|
||||
if len(issued) != 1 {
|
||||
t.Fatalf("expected one route on the bus, got %v", issued)
|
||||
}
|
||||
if port, ok := asPort(issued[0].Values["port"]); !ok || port != 20000 {
|
||||
t.Errorf("the bus carries a port nothing listens on: %v", issued[0].Values["port"])
|
||||
}
|
||||
|
||||
// A module that receives nothing is issued nothing to receive.
|
||||
if _, any := composed.Received["gitea"]; any {
|
||||
t.Errorf("a module that receives nothing was issued something: %v", composed.Received["gitea"])
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user