Compare commits
1
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
24f024dd74 |
@@ -404,7 +404,8 @@ func declarationWith(ctx context.Context, open *stores, node string,
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
return sendable{}, err
|
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.
|
// 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))
|
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
|
// 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.
|
// operators run, and on 2026-10-01 it was the one path that issued none.
|
||||||
var told []string
|
if err := issueMemberships(ctx, open, server, sending); err != nil {
|
||||||
for _, s := range sending {
|
|
||||||
told = append(told, s.node)
|
|
||||||
}
|
|
||||||
if err := issueMemberships(ctx, open, server, told); err != nil {
|
|
||||||
return err
|
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
|
// 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
|
// 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.
|
// 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.
|
// issueMemberships publishes the membership of every module on the machines just sent.
|
||||||
func issueMemberships(ctx context.Context, open *stores, server *link.Server, names []string) error {
|
//
|
||||||
|
// 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)
|
records, err := open.inventory.BusRecords(ctx)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
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.
|
// the push stands, the first failure is named once, and the next push tries again.
|
||||||
issued, failed := 0, 0
|
issued, failed := 0, 0
|
||||||
var first error
|
var first error
|
||||||
for _, node := range names {
|
for _, s := range sent {
|
||||||
|
node := s.node
|
||||||
for _, d := range records.Assigned[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 {
|
if err != nil {
|
||||||
return err
|
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 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 existed: an older host parses the envelope strictly and would refuse the key.
|
||||||
Adoption *adoptionEnvelope
|
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
|
// 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"
|
"sort"
|
||||||
"strings"
|
"strings"
|
||||||
"sync"
|
"sync"
|
||||||
|
"sync/atomic"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
"golang.org/x/crypto/acme"
|
"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
|
// 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.
|
// not merely pointless but the failing order onlyWhatTheMeshSaid exists to prevent.
|
||||||
public map[string]bool
|
public map[string]bool
|
||||||
// inside is the private network's range, where a request must come from to be served a name
|
// inside is where a request must come from to be served a name that is only internal: every
|
||||||
// that is only internal. Set once at start, never replaced with the routes: it is what the
|
// machine's address on the private network, as the mesh issued it in this proxy's membership
|
||||||
// private network is, not what is routed on it.
|
// (novox/hq ADR 0167). Empty until it is issued, and then only the machine itself is inside.
|
||||||
inside sources
|
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 —
|
// sources is who may be served an internal name: the private network's addresses as the mesh
|
||||||
// anything on a machine may call anything on it (novox/hq ADR 0144) — so loopback needs no range.
|
// 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
|
type sources []netip.Prefix
|
||||||
|
|
||||||
// sourcesFrom reads the ranges the mesh wrote, separated by commas or spaces. A range that does not
|
// sourcesOf reads the addresses the mesh issued, each a single address or a range. One that does
|
||||||
// parse is an error, not a range skipped: the proxy would otherwise serve internal names to fewer
|
// not parse is an error, not an entry skipped: the proxy would otherwise serve internal names to
|
||||||
// machines than the mesh said, or start believing a typo.
|
// fewer machines than the mesh said, and say nothing.
|
||||||
func sourcesFrom(text string) (sources, error) {
|
func sourcesOf(mesh []string) (sources, error) {
|
||||||
var out sources
|
var out sources
|
||||||
for _, field := range strings.FieldsFunc(text, func(r rune) bool { return r == ',' || r == ' ' || r == '\n' || r == '\t' }) {
|
for _, entry := range mesh {
|
||||||
prefix, err := netip.ParsePrefix(field)
|
entry = strings.TrimSpace(entry)
|
||||||
if err != nil {
|
if prefix, err := netip.ParsePrefix(entry); err == nil {
|
||||||
return nil, fmt.Errorf("%q is not an address range: %w", field, err)
|
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
|
return out, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// bridgesFrom is the address ranges of the machine's container bridges — the same interfaces the
|
// holds says whether a request from this remote address came from the mesh or the machine itself.
|
||||||
// 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.
|
|
||||||
//
|
//
|
||||||
// **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
|
// 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
|
// 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 {
|
func (s sources) holds(remote string) bool {
|
||||||
host := remote
|
host := remote
|
||||||
if h, _, err := net.SplitHostPort(remote); err == nil {
|
if h, _, err := net.SplitHostPort(remote); err == nil {
|
||||||
@@ -418,13 +386,13 @@ func (t *table) hiddenFrom(host, remote string) bool {
|
|||||||
}
|
}
|
||||||
t.mu.RLock()
|
t.mu.RLock()
|
||||||
defer t.mu.RUnlock()
|
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.
|
// setInside replaces who the mesh is, as the membership said.
|
||||||
func (t *table) setBridges(bridges sources) {
|
func (t *table) setInside(inside sources) {
|
||||||
t.mu.Lock()
|
t.mu.Lock()
|
||||||
t.bridges = bridges
|
t.inside = inside
|
||||||
t.mu.Unlock()
|
t.mu.Unlock()
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -469,19 +437,20 @@ func run() error {
|
|||||||
}
|
}
|
||||||
|
|
||||||
held := newTable()
|
held := newTable()
|
||||||
// Unset means only this machine and its containers are inside, which serves an internal-only
|
// **The bus first, the file until it has spoken** (novox/hq ADR 0167). The membership carries
|
||||||
// name to nobody else — refused rather than served to everyone, which is what the proxy did
|
// the routes and who the mesh is; the file carries the routes alone, so while the proxy reads
|
||||||
// before it knew.
|
// it an internal name is served to this machine and to nobody else — refused, never opened.
|
||||||
inside, err := sourcesFrom(os.Getenv("INTERNAL_SOURCES"))
|
fromBus := &atomic.Bool{}
|
||||||
if err != nil {
|
if credential := strings.TrimSpace(os.Getenv("MESH_BROKER_FILE")); credential != "" {
|
||||||
return fmt.Errorf("INTERNAL_SOURCES: %w", err)
|
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() {
|
read := func() {
|
||||||
|
if fromBus.Load() {
|
||||||
|
return
|
||||||
|
}
|
||||||
routes, public, err := routesFrom(path)
|
routes, public, err := routesFrom(path)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
// Kept serving what it had. A file being rewritten is momentarily unreadable, and
|
// Kept serving what it had. A file being rewritten is momentarily unreadable, and
|
||||||
@@ -491,7 +460,6 @@ func run() error {
|
|||||||
return
|
return
|
||||||
}
|
}
|
||||||
held.set(routes, public)
|
held.set(routes, public)
|
||||||
held.setBridges(theseBridges())
|
|
||||||
log.Printf("serving %d route(s): %s", len(routes), strings.Join(held.names(), ", "))
|
log.Printf("serving %d route(s): %s", len(routes), strings.Join(held.names(), ", "))
|
||||||
}
|
}
|
||||||
read()
|
read()
|
||||||
@@ -886,10 +854,16 @@ func routesFrom(path string) (map[string][]rule, map[string]bool, error) {
|
|||||||
if err := json.Unmarshal(raw, &said); err != nil {
|
if err := json.Unmarshal(raw, &said); err != nil {
|
||||||
return nil, nil, err
|
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{}
|
out := map[string][]rule{}
|
||||||
public := map[string]bool{}
|
public := map[string]bool{}
|
||||||
for _, c := range said.Given {
|
for _, c := range contributions {
|
||||||
name, _ := c.Values["name"].(string)
|
name, _ := c.Values["name"].(string)
|
||||||
name = strings.TrimSpace(name)
|
name = strings.TrimSpace(name)
|
||||||
internal, _ := c.Values["internal-name"].(string)
|
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)
|
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.
|
// asWhole is any whole number the mesh wrote, whatever its magnitude.
|
||||||
|
|||||||
@@ -2,6 +2,7 @@ package main
|
|||||||
|
|
||||||
import (
|
import (
|
||||||
"crypto/tls"
|
"crypto/tls"
|
||||||
|
"encoding/json"
|
||||||
"fmt"
|
"fmt"
|
||||||
"io"
|
"io"
|
||||||
"net"
|
"net"
|
||||||
@@ -10,11 +11,13 @@ import (
|
|||||||
"net/url"
|
"net/url"
|
||||||
"strings"
|
"strings"
|
||||||
"testing"
|
"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
|
// 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.
|
// internal-only name to it, with the mesh's machines as the membership would issue them.
|
||||||
func behind(t *testing.T, inside string) *table {
|
func behind(t *testing.T, mesh ...string) *table {
|
||||||
t.Helper()
|
t.Helper()
|
||||||
workload := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
workload := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||||
io.WriteString(w, "the workload")
|
io.WriteString(w, "the workload")
|
||||||
@@ -33,10 +36,11 @@ func behind(t *testing.T, inside string) *table {
|
|||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
held := newTable()
|
held := newTable()
|
||||||
held.inside, err = sourcesFrom(inside)
|
inside, err := sourcesOf(mesh)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
|
held.setInside(inside)
|
||||||
held.set(routes, public)
|
held.set(routes, public)
|
||||||
return held
|
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
|
// 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.
|
// being internal kept nobody out: a request from the internet only had to carry it.
|
||||||
func TestAnInternalOnlyNameIsServedOnlyInsideThePrivateNetwork(t *testing.T) {
|
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 ||
|
if code, body := askFrom(held, "admin.anchor.internal", "10.10.0.7:51000"); code != http.StatusOK ||
|
||||||
body != "the workload" {
|
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
|
// 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.
|
// name stays one thing whichever route it came from.
|
||||||
func TestAnInternalAliasOfAPublicRouteIsServedInsideOnly(t *testing.T) {
|
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 {
|
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)
|
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
|
// Before a membership has said who the mesh is, only the machine itself is inside — refused to
|
||||||
// (novox/hq ADR 0144): inside, though it is neither loopback nor the mesh.
|
// everyone else, never served to everyone.
|
||||||
func TestAContainerOnThisMachineIsInside(t *testing.T) {
|
func TestUntilTheMeshIsIssuedAnInternalOnlyNameIsServedToTheMachineAlone(t *testing.T) {
|
||||||
held := behind(t, "10.10.0.0/24")
|
held := behind(t)
|
||||||
_, 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, "")
|
|
||||||
if code, _ := askFrom(held, "admin.anchor.internal", "10.10.0.7:51000"); code != http.StatusNotFound {
|
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)
|
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,
|
// 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.
|
// and serving it would answer the question the routing refuses to.
|
||||||
func TestTheHandshakeRefusesAnInternalOnlyNameToAnOutsider(t *testing.T) {
|
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{}
|
served := &tls.Certificate{}
|
||||||
pick := certificateFor(held, func(*tls.ClientHelloInfo) (*tls.Certificate, error) { return served, nil }, nil)
|
pick := certificateFor(held, func(*tls.ClientHelloInfo) (*tls.Certificate, error) { return served, nil }, nil)
|
||||||
hello := func(name, remote string) *tls.ClientHelloInfo {
|
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
|
// The mesh is issued as machines' addresses; a range is read as well. One that does not parse is
|
||||||
// than the mesh said, or to a typo.
|
// refused rather than skipped, so a typo never quietly narrows or widens who is inside.
|
||||||
func TestAPrivateNetworkThatDoesNotParseIsRefused(t *testing.T) {
|
func TestTheMeshIsReadAsAddressesAndRanges(t *testing.T) {
|
||||||
if _, err := sourcesFrom("10.10.0.0/24, not-a-range"); err == nil {
|
if _, err := sourcesOf([]string{"10.10.0.1", "not-an-address"}); err == nil {
|
||||||
t.Error("a range that does not parse was accepted")
|
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 {
|
if err != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
for remote, want := range map[string]bool{
|
for remote, want := range map[string]bool{
|
||||||
"10.10.0.200:1": true,
|
"10.10.0.1:1": true,
|
||||||
"[::ffff:10.10.0.3]:1": true,
|
"[::ffff:10.10.0.1]:1": true,
|
||||||
"[fd00::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,
|
"192.168.1.10:1": false,
|
||||||
"not-an-address": false,
|
"not-an-address": false,
|
||||||
} {
|
} {
|
||||||
if inside.holds(remote) != want {
|
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
|
package broker
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"encoding/json"
|
||||||
"sort"
|
"sort"
|
||||||
"strings"
|
"strings"
|
||||||
)
|
)
|
||||||
@@ -33,6 +34,17 @@ type Membership struct {
|
|||||||
Reaches map[string][]string `json:"reaches,omitempty"`
|
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 is where this instance answers what it serves — the runtime's one verb of its own.
|
||||||
Tools string `json:"tools"`
|
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.
|
// 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:
|
// 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 —
|
// a module's name may itself contain a dot. What the mesh adds of its own — an opening, the guard —
|
||||||
// has no owner.
|
// 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 {
|
type Composed struct {
|
||||||
Resources []map[string]any
|
Resources []map[string]any
|
||||||
Owner map[string]string
|
Owner map[string]string
|
||||||
|
Received map[string]map[string][]Contribution
|
||||||
}
|
}
|
||||||
|
|
||||||
// Compose is Declaration with the owner of every resource said.
|
// Compose is Declaration with the owner of every resource said.
|
||||||
func (r Resolution) Compose(with Rendering) (Composed, error) {
|
func (r Resolution) Compose(with Rendering) (Composed, error) {
|
||||||
owner := map[string]string{}
|
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 {
|
if err != nil {
|
||||||
return Composed{}, err
|
return Composed{}, err
|
||||||
}
|
}
|
||||||
@@ -261,7 +267,7 @@ func (r Resolution) Compose(with Rendering) (Composed, error) {
|
|||||||
"sealed": with.BusMembership, "mode": "0600",
|
"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
|
// 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"
|
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,
|
// 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
|
// 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
|
// 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
|
return nil, err
|
||||||
}
|
}
|
||||||
first = append(first, file)
|
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 {
|
if m.Keeps != "" && with.Kept != nil {
|
||||||
file, err := keptFile(m.Keeps, with.Kept)
|
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