Compare commits
18
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
e7da39de57 | ||
|
|
e6ddc59cde | ||
|
|
6c5dfd0c25 | ||
|
|
775df79893 | ||
|
|
3fbf658c16 | ||
|
|
bc31745607 | ||
|
|
e5c2eb20f2 | ||
|
|
05fb7fb5eb | ||
|
|
2c1733de8d | ||
|
|
e51c94dcb5 | ||
|
|
07c07902ff | ||
|
|
864cdea4c6 | ||
|
|
6e810907b2 | ||
|
|
2b20a12c4a | ||
|
|
9c83dacfce | ||
|
|
64ba053f3b | ||
|
|
96416bd8a7 | ||
|
|
aaad02fd38 |
@@ -27,8 +27,18 @@ build:
|
|||||||
IMAGE ?= mesh-controller:$(VERSION)
|
IMAGE ?= mesh-controller:$(VERSION)
|
||||||
DEV_TAG ?= mesh-controller:development
|
DEV_TAG ?= mesh-controller:development
|
||||||
|
|
||||||
|
# The base the module declares, read from the manifest rather than written here twice.
|
||||||
|
#
|
||||||
|
# **`make image` was broken and stayed broken**, because the Dockerfile's fallback base was a Go
|
||||||
|
# older than go.mod asks for: every build died at `go mod download` with "go.mod requires go >=
|
||||||
|
# 1.26.0", and the pipeline never saw it because the pipeline passes the declared base in. Anybody
|
||||||
|
# building the image by hand hit it and had to find the digest themselves (novox/hq 04-ISSUES/146,
|
||||||
|
# what it cost).
|
||||||
|
GO_BASE ?= $(shell python3 -c "import json;print(next(o['image'] for o in json.load(open('module.json'))['build']['on'] if o['arg']=='GO_BASE'))" 2>/dev/null)
|
||||||
|
|
||||||
image:
|
image:
|
||||||
docker build --build-arg VERSION=$(VERSION) -t $(IMAGE) -t $(DEV_TAG) .
|
@test -n "$(GO_BASE)" || { echo "module.json declares no GO_BASE; pass GO_BASE=<image> or fix the manifest"; exit 1; }
|
||||||
|
docker build --build-arg GO_BASE=$(GO_BASE) --build-arg VERSION=$(VERSION) -t $(IMAGE) -t $(DEV_TAG) .
|
||||||
@echo
|
@echo
|
||||||
@docker image inspect $(IMAGE) --format 'built {{.RepoTags}} {{.Size}} bytes'
|
@docker image inspect $(IMAGE) --format 'built {{.RepoTags}} {{.Size}} bytes'
|
||||||
|
|
||||||
@@ -38,7 +48,8 @@ BUILDER_IMAGE ?= mesh-builder:$(VERSION)
|
|||||||
BUILDER_DEV_TAG ?= mesh-builder:development
|
BUILDER_DEV_TAG ?= mesh-builder:development
|
||||||
|
|
||||||
builder-image:
|
builder-image:
|
||||||
docker build -f cmd/mesh-builder/Dockerfile -t $(BUILDER_IMAGE) -t $(BUILDER_DEV_TAG) .
|
@test -n "$(GO_BASE)" || { echo "module.json declares no GO_BASE; pass GO_BASE=<image> or fix the manifest"; exit 1; }
|
||||||
|
docker build --build-arg GO_BASE=$(GO_BASE) -f cmd/mesh-builder/Dockerfile -t $(BUILDER_IMAGE) -t $(BUILDER_DEV_TAG) .
|
||||||
@echo
|
@echo
|
||||||
@docker image inspect $(BUILDER_IMAGE) --format 'built {{.RepoTags}} {{.Size}} bytes'
|
@docker image inspect $(BUILDER_IMAGE) --format 'built {{.RepoTags}} {{.Size}} bytes'
|
||||||
|
|
||||||
|
|||||||
@@ -0,0 +1,258 @@
|
|||||||
|
package main
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"crypto/rand"
|
||||||
|
"crypto/rsa"
|
||||||
|
"crypto/x509"
|
||||||
|
"crypto/x509/pkix"
|
||||||
|
"encoding/pem"
|
||||||
|
"errors"
|
||||||
|
"fmt"
|
||||||
|
"math/big"
|
||||||
|
"net"
|
||||||
|
"os"
|
||||||
|
"path/filepath"
|
||||||
|
"strings"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"github.com/novox/mesh-controller/internal/broker"
|
||||||
|
)
|
||||||
|
|
||||||
|
// The bus's own certificate, made by the mesh rather than borrowed from an image.
|
||||||
|
//
|
||||||
|
// **The foundation asked a third-party image for a tool it never said must be there** (novox/hq
|
||||||
|
// 04-ISSUES/146). The bootstrap made this certificate by running `openssl` inside the broker's
|
||||||
|
// image, which worked while the broker was one that happened to carry it and stopped the day the
|
||||||
|
// bus changed: the new one has a shell and no openssl, so the step exited 127 and no mesh could be
|
||||||
|
// raised. Substituting another image the bundle names does not help — none of them carry it
|
||||||
|
// either.
|
||||||
|
//
|
||||||
|
// So the program that needs a certificate makes one. It is the mesh's own binary, already on the
|
||||||
|
// machine at this point in the bootstrap (the schema step ran it), and it needs nothing from the
|
||||||
|
// image it writes into but a mounted directory.
|
||||||
|
//
|
||||||
|
// **Self-signed, and that is the design** — a host pins this server's exact certificate and
|
||||||
|
// authenticates with a password (novox/hq ADR 0004). There is no authority above it to ask, and at
|
||||||
|
// this moment in a bootstrap there is no mesh to ask one of.
|
||||||
|
//
|
||||||
|
// Idempotent, because the step is applied again on every reconcile and a second certificate would
|
||||||
|
// be one the hosts that pinned the first no longer believe.
|
||||||
|
|
||||||
|
// busCertificateNames is what the bus is reached by: the container name on a mesh network, and the
|
||||||
|
// loopback address the machine's own foundation dials.
|
||||||
|
var busCertificateNames = []string{"mesh-broker"}
|
||||||
|
|
||||||
|
const busCertificateLife = 10 * 365 * 24 * time.Hour
|
||||||
|
|
||||||
|
// busCertificate makes the bus's certificate in a directory, or says whether one is there.
|
||||||
|
//
|
||||||
|
// broker certificate --into /tls make it if it is not there
|
||||||
|
// broker certificate --check --into /tls exit non-zero unless a usable pair is
|
||||||
|
func busCertificate(args []string) error {
|
||||||
|
into, check := "", false
|
||||||
|
for i := 0; i < len(args); i++ {
|
||||||
|
switch args[i] {
|
||||||
|
case "--check":
|
||||||
|
check = true
|
||||||
|
case "--into":
|
||||||
|
if i+1 >= len(args) {
|
||||||
|
return errors.New("--into needs a directory")
|
||||||
|
}
|
||||||
|
into = args[i+1]
|
||||||
|
i++
|
||||||
|
default:
|
||||||
|
return fmt.Errorf("broker certificate [--check] --into <directory>: %q", args[i])
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if into == "" {
|
||||||
|
return errors.New("broker certificate [--check] --into <directory>")
|
||||||
|
}
|
||||||
|
crt, key := filepath.Join(into, "tls.crt"), filepath.Join(into, "tls.key")
|
||||||
|
|
||||||
|
if usable, err := busCertificateUsable(crt, key); err != nil {
|
||||||
|
return err
|
||||||
|
} else if usable {
|
||||||
|
fmt.Printf("the bus already has a certificate at %s, and it was left alone\n", crt)
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
if check {
|
||||||
|
// Said as a failure, because that is what the caller asked: a bootstrap's verify runs
|
||||||
|
// this and a false answer is what makes the step run.
|
||||||
|
return fmt.Errorf("no usable certificate and key at %s", into)
|
||||||
|
}
|
||||||
|
return writeBusCertificate(crt, key)
|
||||||
|
}
|
||||||
|
|
||||||
|
// busCertificateUsable says whether a certificate and its key are both there and parse.
|
||||||
|
//
|
||||||
|
// Both, and parsed rather than stat'ed: a half-written pair is the state a bootstrap interrupted
|
||||||
|
// between the two files leaves behind, and a step that treated it as done would hand the server a
|
||||||
|
// certificate with no key and report success.
|
||||||
|
func busCertificateUsable(crt, key string) (bool, error) {
|
||||||
|
certPEM, err := os.ReadFile(crt)
|
||||||
|
if errors.Is(err, os.ErrNotExist) {
|
||||||
|
return false, nil
|
||||||
|
}
|
||||||
|
if err != nil {
|
||||||
|
return false, err
|
||||||
|
}
|
||||||
|
keyPEM, err := os.ReadFile(key)
|
||||||
|
if errors.Is(err, os.ErrNotExist) {
|
||||||
|
return false, nil
|
||||||
|
}
|
||||||
|
if err != nil {
|
||||||
|
return false, err
|
||||||
|
}
|
||||||
|
if _, err := tlsPairParses(certPEM, keyPEM); err != nil {
|
||||||
|
return false, nil
|
||||||
|
}
|
||||||
|
return true, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func tlsPairParses(certPEM, keyPEM []byte) (*x509.Certificate, error) {
|
||||||
|
block, _ := pem.Decode(certPEM)
|
||||||
|
if block == nil || block.Type != "CERTIFICATE" {
|
||||||
|
return nil, errors.New("not a certificate")
|
||||||
|
}
|
||||||
|
certificate, err := x509.ParseCertificate(block.Bytes)
|
||||||
|
if err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
keyBlock, _ := pem.Decode(keyPEM)
|
||||||
|
if keyBlock == nil {
|
||||||
|
return nil, errors.New("not a key")
|
||||||
|
}
|
||||||
|
if _, err := x509.ParsePKCS8PrivateKey(keyBlock.Bytes); err != nil {
|
||||||
|
if _, err := x509.ParsePKCS1PrivateKey(keyBlock.Bytes); err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return certificate, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func writeBusCertificate(crt, key string) error {
|
||||||
|
private, err := rsa.GenerateKey(rand.Reader, 2048)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
serial, err := rand.Int(rand.Reader, new(big.Int).Lsh(big.NewInt(1), 128))
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
template := &x509.Certificate{
|
||||||
|
SerialNumber: serial,
|
||||||
|
Subject: pkix.Name{CommonName: busCertificateNames[0]},
|
||||||
|
DNSNames: busCertificateNames,
|
||||||
|
IPAddresses: []net.IP{net.ParseIP("127.0.0.1")},
|
||||||
|
NotBefore: time.Now().Add(-time.Hour),
|
||||||
|
NotAfter: time.Now().Add(busCertificateLife),
|
||||||
|
KeyUsage: x509.KeyUsageDigitalSignature | x509.KeyUsageKeyEncipherment,
|
||||||
|
ExtKeyUsage: []x509.ExtKeyUsage{x509.ExtKeyUsageServerAuth},
|
||||||
|
BasicConstraintsValid: true,
|
||||||
|
}
|
||||||
|
der, err := x509.CreateCertificate(rand.Reader, template, template, &private.PublicKey, private)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
pkcs8, err := x509.MarshalPKCS8PrivateKey(private)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
|
||||||
|
// **The key first, and only then the certificate**, so the pair a reader finds is never a
|
||||||
|
// certificate whose key has not been written yet — the one order in which an interruption
|
||||||
|
// leaves something that looks finished (novox/hq 04-ISSUES/014, a key present and unusable).
|
||||||
|
if err := os.WriteFile(key, pem.EncodeToMemory(&pem.Block{Type: "PRIVATE KEY", Bytes: pkcs8}), 0o600); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
if err := os.WriteFile(crt, pem.EncodeToMemory(&pem.Block{Type: "CERTIFICATE", Bytes: der}), 0o644); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
fmt.Printf("made the bus a certificate for %v, valid until %s\n %s\n %s\n",
|
||||||
|
busCertificateNames, template.NotAfter.Format(time.RFC3339), crt, key)
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// busAccounts writes the mesh's composed user list to a file.
|
||||||
|
//
|
||||||
|
// **For genesis, where no declaration can deliver it** (novox/hq 04-ISSUES/146). Everywhere else
|
||||||
|
// the list reaches the machine running the bus as a resource of the module that holds it — which
|
||||||
|
// requires that machine to be an enrolled node, and at genesis it is not: the first node cannot
|
||||||
|
// enrol because the account it would enrol with cannot be composed onto a bus it has no declaration
|
||||||
|
// for. The installer breaks that circle by placing the file itself, once, and the module takes the
|
||||||
|
// file over from its first push.
|
||||||
|
//
|
||||||
|
// The same composition, not a second one: this asks the store for the same records and renders them
|
||||||
|
// with the same composer the declaration uses. A genesis that hand-wrote an account would be a
|
||||||
|
// second statement of who may say what, able to disagree with the first.
|
||||||
|
//
|
||||||
|
// **It writes to standard output unless told a file**, and that is the point: the control plane
|
||||||
|
// composes and says what it composed, and whoever is raising the machine puts it where that
|
||||||
|
// machine's bus reads it. A control plane that wrote into the bus's own directory would have to
|
||||||
|
// know where that is and how to make the server re-read it — which is the module's knowledge, and
|
||||||
|
// the module is what takes this over on the first push.
|
||||||
|
//
|
||||||
|
// broker accounts > /var/lib/mesh-bus-conf/accounts.conf
|
||||||
|
func busAccounts(ctx context.Context, args []string) error {
|
||||||
|
into := ""
|
||||||
|
for i := 0; i < len(args); i++ {
|
||||||
|
switch args[i] {
|
||||||
|
case "--into":
|
||||||
|
if i+1 >= len(args) {
|
||||||
|
return errors.New("--into needs a file")
|
||||||
|
}
|
||||||
|
into = args[i+1]
|
||||||
|
i++
|
||||||
|
default:
|
||||||
|
return fmt.Errorf("broker accounts --into <file>: %q", args[i])
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
open, err := openStores(ctx)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
defer open.Close()
|
||||||
|
|
||||||
|
records, err := open.inventory.BusRecords(ctx)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
users, err := broker.Users(records)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
kept, err := open.inventory.BusUsers(ctx)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
hashes := make(map[string]string, len(kept))
|
||||||
|
for name, u := range kept {
|
||||||
|
hashes[name] = u.PasswordHash
|
||||||
|
}
|
||||||
|
filled, missing := broker.WithPasswords(users, hashes)
|
||||||
|
if len(missing) > 0 {
|
||||||
|
// To standard error, always: the composed file may be going to standard output, and a
|
||||||
|
// remark in the middle of it is a configuration the server refuses to parse.
|
||||||
|
fmt.Fprintf(os.Stderr, "leaving out %d user(s) the mesh has minted no credential for: %s\n",
|
||||||
|
len(missing), strings.Join(missing, ", "))
|
||||||
|
}
|
||||||
|
if len(filled) == 0 {
|
||||||
|
return errors.New("not one user has a credential, so this list would refuse every " +
|
||||||
|
"connection in the mesh")
|
||||||
|
}
|
||||||
|
accounts, err := broker.ComposeAccounts(filled)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
if into == "" {
|
||||||
|
fmt.Print(accounts)
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
if err := os.WriteFile(into, []byte(accounts), 0o600); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
fmt.Printf("wrote %d user(s) to %s\n", len(filled), into)
|
||||||
|
return nil
|
||||||
|
}
|
||||||
@@ -0,0 +1,121 @@
|
|||||||
|
package main
|
||||||
|
|
||||||
|
import (
|
||||||
|
"crypto/tls"
|
||||||
|
"crypto/x509"
|
||||||
|
"os"
|
||||||
|
"path/filepath"
|
||||||
|
"strings"
|
||||||
|
"testing"
|
||||||
|
)
|
||||||
|
|
||||||
|
// novox/hq 04-ISSUES/146. The bootstrap could not make the bus a certificate: it asked an image for
|
||||||
|
// `openssl` and the image it asks has none. What replaces it is this command, so what is checked is
|
||||||
|
// what the bootstrap needs from it — a pair a TLS server can actually load, made once and only once.
|
||||||
|
|
||||||
|
func TestTheBusCertificateLoadsAsAServersWould(t *testing.T) {
|
||||||
|
into := t.TempDir()
|
||||||
|
if err := busCertificate([]string{"--into", into}); err != nil {
|
||||||
|
t.Fatalf("the bus could not be given a certificate: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
// The check a cheaper test would not make. The key was present and valid and the server could
|
||||||
|
// not start, once, because nothing loaded the pair the way a server loads it
|
||||||
|
// (novox/hq 04-ISSUES/014).
|
||||||
|
pair, err := tls.LoadX509KeyPair(filepath.Join(into, "tls.crt"), filepath.Join(into, "tls.key"))
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("a TLS server cannot load what was written: %v", err)
|
||||||
|
}
|
||||||
|
leaf := pair.Leaf
|
||||||
|
if leaf == nil {
|
||||||
|
if leaf, err = x509.ParseCertificate(pair.Certificate[0]); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if err := leaf.VerifyHostname("mesh-broker"); err != nil {
|
||||||
|
t.Errorf("the certificate is not for the name the bus is reached by: %v", err)
|
||||||
|
}
|
||||||
|
if len(leaf.IPAddresses) == 0 || leaf.IPAddresses[0].String() != "127.0.0.1" {
|
||||||
|
t.Errorf("the certificate does not cover the loopback address the foundation dials: %v", leaf.IPAddresses)
|
||||||
|
}
|
||||||
|
|
||||||
|
// The key is not readable by anything else on the machine; the certificate is public and is.
|
||||||
|
key, err := os.Stat(filepath.Join(into, "tls.key"))
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if key.Mode().Perm() != 0o600 {
|
||||||
|
t.Errorf("the key is %v", key.Mode().Perm())
|
||||||
|
}
|
||||||
|
crt, err := os.Stat(filepath.Join(into, "tls.crt"))
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if crt.Mode().Perm() != 0o644 {
|
||||||
|
t.Errorf("the certificate is %v, which the server runs as another user cannot read", crt.Mode().Perm())
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// **Made once.** The step is applied again on every reconcile, and a second certificate is one the
|
||||||
|
// hosts that pinned the first no longer believe (novox/hq ADR 0004).
|
||||||
|
func TestTheBusCertificateIsMadeOnce(t *testing.T) {
|
||||||
|
into := t.TempDir()
|
||||||
|
if err := busCertificate([]string{"--into", into}); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
first, err := os.ReadFile(filepath.Join(into, "tls.crt"))
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if err := busCertificate([]string{"--into", into}); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
again, err := os.ReadFile(filepath.Join(into, "tls.crt"))
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if string(first) != string(again) {
|
||||||
|
t.Fatal("running it twice replaced the certificate every host had pinned")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// The verify half: false before, true after, which is what makes the bootstrap run the step at all.
|
||||||
|
func TestTheCheckIsFalseUntilThereIsAPair(t *testing.T) {
|
||||||
|
into := t.TempDir()
|
||||||
|
if err := busCertificate([]string{"--check", "--into", into}); err == nil {
|
||||||
|
t.Fatal("an empty directory reported a usable certificate")
|
||||||
|
}
|
||||||
|
if err := busCertificate([]string{"--into", into}); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if err := busCertificate([]string{"--check", "--into", into}); err != nil {
|
||||||
|
t.Fatalf("the certificate it just made does not satisfy its own check: %v", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// A half-written pair is not a pair. An interrupted bootstrap leaves exactly this, and a step that
|
||||||
|
// called it done would hand the server a certificate with no key and report success.
|
||||||
|
func TestACertificateWithoutItsKeyIsNotUsable(t *testing.T) {
|
||||||
|
into := t.TempDir()
|
||||||
|
if err := busCertificate([]string{"--into", into}); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if err := os.Remove(filepath.Join(into, "tls.key")); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if err := busCertificate([]string{"--check", "--into", into}); err == nil {
|
||||||
|
t.Fatal("a certificate with no key passed the check")
|
||||||
|
}
|
||||||
|
if err := busCertificate([]string{"--into", into}); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if _, err := tls.LoadX509KeyPair(filepath.Join(into, "tls.crt"), filepath.Join(into, "tls.key")); err != nil {
|
||||||
|
t.Fatalf("it did not replace the unusable pair: %v", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestWhereToWriteIsRequired(t *testing.T) {
|
||||||
|
if err := busCertificate(nil); err == nil || !strings.Contains(err.Error(), "--into") {
|
||||||
|
t.Fatalf("it did not ask where to write: %v", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -88,7 +88,7 @@ func run() error {
|
|||||||
case "identity":
|
case "identity":
|
||||||
return identityCommand(ctx, args[1:])
|
return identityCommand(ctx, args[1:])
|
||||||
case "broker":
|
case "broker":
|
||||||
return brokerCommand(args[1:])
|
return brokerCommand(ctx, args[1:])
|
||||||
case "serve":
|
case "serve":
|
||||||
return serve(ctx)
|
return serve(ctx)
|
||||||
case "upgrade":
|
case "upgrade":
|
||||||
|
|||||||
@@ -346,9 +346,16 @@ func whoResolves(ctx context.Context, open *stores, requirement string) (
|
|||||||
refused := map[string]string{}
|
refused := map[string]string{}
|
||||||
for _, n := range nodes {
|
for _, n := range nodes {
|
||||||
plan, _, err := planFor(ctx, open, n.Name)
|
plan, _, err := planFor(ctx, open, n.Name)
|
||||||
if err != nil {
|
switch {
|
||||||
|
case unresolvable(err):
|
||||||
refused[n.Name] = err.Error()
|
refused[n.Name] = err.Error()
|
||||||
continue
|
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 _, m := range plan.Modules {
|
||||||
for _, offered := range m.Offers() {
|
for _, offered := range m.Offers() {
|
||||||
|
|||||||
@@ -363,9 +363,15 @@ func identityCommand(ctx context.Context, args []string) error {
|
|||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func brokerCommand(args []string) error {
|
func brokerCommand(ctx context.Context, args []string) error {
|
||||||
|
if len(args) > 0 && args[0] == "certificate" {
|
||||||
|
return busCertificate(args[1:])
|
||||||
|
}
|
||||||
|
if len(args) > 0 && args[0] == "accounts" {
|
||||||
|
return busAccounts(ctx, args[1:])
|
||||||
|
}
|
||||||
if len(args) == 0 || args[0] != "show" {
|
if len(args) == 0 || args[0] != "show" {
|
||||||
return errors.New("broker show")
|
return errors.New("broker show | broker certificate [--check] --into <directory> | broker accounts --into <file>")
|
||||||
}
|
}
|
||||||
known, err := broker.FromEnvironment()
|
known, err := broker.FromEnvironment()
|
||||||
if errors.Is(err, broker.ErrNotConfigured) {
|
if errors.Is(err, broker.ErrNotConfigured) {
|
||||||
|
|||||||
@@ -25,7 +25,36 @@ import (
|
|||||||
// cheapest next step. That is how novox/hq ADR 0001 records `hal/sdk` reaching 34,636:
|
// 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.
|
// 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.
|
// 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) {
|
func planFor(ctx context.Context, open *stores, nodeName string) (catalogue.Resolution, catalogue.SettingsBy, error) {
|
||||||
inv := open.inventory
|
inv := open.inventory
|
||||||
shelf, err := inv.Catalogue(ctx)
|
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,
|
At: onNetwork[nodeName], PublicDomain: publicDomain,
|
||||||
Account: who.Account, AccountHome: who.AccountHome}, world)
|
Account: who.Account, AccountHome: who.AccountHome}, world)
|
||||||
if err != nil {
|
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
|
// 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 {
|
if len(stray) > 0 {
|
||||||
// Somebody set something that reaches no file. Said here rather than discovered by the
|
// 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.
|
// 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
|
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
|
// 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
|
// "route" means. A node that does not resolve is skipped, so one machine's broken set does not cost
|
||||||
// the rest their names.
|
// 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) {
|
func routeNamesInTheMesh(ctx context.Context, open *stores) (map[string]string, error) {
|
||||||
inv := open.inventory
|
inv := open.inventory
|
||||||
places, err := inv.Overlays(ctx)
|
places, err := inv.Overlays(ctx)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
return nil, fmt.Errorf("where the machines are cannot be read: %w", err)
|
||||||
}
|
}
|
||||||
address := map[string]string{}
|
address := map[string]string{}
|
||||||
for _, p := range places {
|
for _, p := range places {
|
||||||
@@ -686,14 +728,22 @@ func routeNamesInTheMesh(ctx context.Context, open *stores) (map[string]string,
|
|||||||
|
|
||||||
nodes, err := inv.Nodes(ctx)
|
nodes, err := inv.Nodes(ctx)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
return nil, fmt.Errorf("which machines the mesh has cannot be read: %w", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
out := map[string]string{}
|
out := map[string]string{}
|
||||||
for _, n := range nodes {
|
for _, n := range nodes {
|
||||||
plan, settings, err := planFor(ctx, open, n.Name)
|
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
|
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 _, m := range plan.Modules {
|
||||||
for to := range m.Contributes {
|
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))
|
out := make([]catalogue.Grant, 0, len(issued))
|
||||||
for _, s := range issued {
|
for _, s := range issued {
|
||||||
plan, settings, err := planFor(ctx, open, s.Consumer)
|
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
|
// 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
|
// 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.
|
// run would have the provider create a user nothing uses.
|
||||||
continue
|
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)
|
values, asks, err := plan.ContributionsFrom(s.Name, s.ConsumerModule, settings)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
|||||||
@@ -717,6 +717,12 @@ func raiseTheBus(ctx context.Context, inv *inventory.Inventory, address string)
|
|||||||
broker.BareAddress(address), err)
|
broker.BareAddress(address), err)
|
||||||
}
|
}
|
||||||
defer js.Close()
|
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
|
// **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
|
// 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))
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -62,6 +62,22 @@ func TestTheInstallersFirstUserListIsWhatTheControllerWouldCompose(t *testing.T)
|
|||||||
|
|
||||||
// theCarriedAccounts is the accounts file the installer's template writes at genesis.
|
// theCarriedAccounts is the accounts file the installer's template writes at genesis.
|
||||||
func theCarriedAccounts(t *testing.T) string {
|
func theCarriedAccounts(t *testing.T) string {
|
||||||
|
t.Helper()
|
||||||
|
for _, r := range theTemplate(t) {
|
||||||
|
if r["id"] == "bus-accounts" {
|
||||||
|
content, _ := r["content"].(string)
|
||||||
|
if content == "" {
|
||||||
|
t.Fatal("the template's accounts file is empty, so the bus would refuse every connection")
|
||||||
|
}
|
||||||
|
return content
|
||||||
|
}
|
||||||
|
}
|
||||||
|
t.Fatal("the template carries no accounts file, so a mesh raised from it has a bus nobody may use")
|
||||||
|
return ""
|
||||||
|
}
|
||||||
|
|
||||||
|
// theTemplate is the installer's bundle, as resources.
|
||||||
|
func theTemplate(t *testing.T) []map[string]any {
|
||||||
t.Helper()
|
t.Helper()
|
||||||
path := filepath.Join("..", "..", "..", "mesh-host", "examples", "foundation-first-node-nats.lock")
|
path := filepath.Join("..", "..", "..", "mesh-host", "examples", "foundation-first-node-nats.lock")
|
||||||
raw, err := os.ReadFile(path)
|
raw, err := os.ReadFile(path)
|
||||||
@@ -81,17 +97,7 @@ func theCarriedAccounts(t *testing.T) string {
|
|||||||
if err := json.Unmarshal([]byte(strings.Join(lines, "\n")), &bundle); err != nil {
|
if err := json.Unmarshal([]byte(strings.Join(lines, "\n")), &bundle); err != nil {
|
||||||
t.Fatalf("the template is not readable: %v", err)
|
t.Fatalf("the template is not readable: %v", err)
|
||||||
}
|
}
|
||||||
for _, r := range bundle.Resources {
|
return bundle.Resources
|
||||||
if r["id"] == "bus-accounts" {
|
|
||||||
content, _ := r["content"].(string)
|
|
||||||
if content == "" {
|
|
||||||
t.Fatal("the template's accounts file is empty, so the bus would refuse every connection")
|
|
||||||
}
|
|
||||||
return content
|
|
||||||
}
|
|
||||||
}
|
|
||||||
t.Fatal("the template carries no accounts file, so a mesh raised from it has a bus nobody may use")
|
|
||||||
return ""
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// subjectsIn reads one allow-list out of a composed accounts file.
|
// subjectsIn reads one allow-list out of a composed accounts file.
|
||||||
|
|||||||
@@ -26,6 +26,16 @@ import (
|
|||||||
type JetStream struct {
|
type JetStream struct {
|
||||||
conn *nats.Conn
|
conn *nats.Conn
|
||||||
js nats.JetStreamContext
|
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.
|
// Dial connects and returns the controller's JetStream handle.
|
||||||
@@ -184,12 +194,60 @@ func (j *JetStream) EnsureConsumer(c Consumer) error {
|
|||||||
// without the other is refused by the server with a message that does not say which half is
|
// without the other is refused by the server with a message that does not say which half is
|
||||||
// missing.
|
// missing.
|
||||||
if c.Queue != "" || c.Push {
|
if c.Queue != "" || c.Push {
|
||||||
want.DeliverSubject = "_DELIVER." + c.Name
|
// **Per consumer, which means per stream as well as per name** (novox/hq 04-ISSUES/146).
|
||||||
|
// A push consumer delivers onto an ordinary subject, and everything subscribed to that
|
||||||
|
// subject gets a copy. The controller holds a consumer called `controller` on CONTROL and
|
||||||
|
// another called `controller` on EVENTS, and both were given `_DELIVER.controller` — so the
|
||||||
|
// one process, holding both subscriptions, acted on every message twice. It enrolled a
|
||||||
|
// joining machine twice from one request, minting a second credential that replaced the one
|
||||||
|
// the machine had just been given; the same doubling applied to every report and every
|
||||||
|
// event the controller follows.
|
||||||
|
//
|
||||||
|
// The stream is in the name because the pair is what identifies a consumer — the server
|
||||||
|
// scopes a durable's name to its stream, and this subject is the only place that scoping
|
||||||
|
// was dropped. Already within what the controller may subscribe (`_DELIVER.controller.>`),
|
||||||
|
// so no permission moves.
|
||||||
|
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:
|
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 {
|
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 fmt.Errorf("bringing consumer %s on %s to match: %w", c.Name, c.Stream, err)
|
||||||
}
|
}
|
||||||
return nil
|
return nil
|
||||||
|
|||||||
@@ -269,7 +269,13 @@ func PermissionsFor(p Principal) (Permissions, error) {
|
|||||||
"mesh.control." + p.Node + ".>",
|
"mesh.control." + p.Node + ".>",
|
||||||
"$JS.API.CONSUMER.INFO.NODES." + p.Node,
|
"$JS.API.CONSUMER.INFO.NODES." + p.Node,
|
||||||
}
|
}
|
||||||
sub = []string{"mesh.node." + p.Node + ".declare", "_DELIVER." + p.Node}
|
// The deliver subject carries the stream as well as the consumer's name, so what a
|
||||||
|
// subscriber is permitted has to carry it too (novox/hq 04-ISSUES/146). The bare name
|
||||||
|
// stays: an existing consumer keeps delivering where it always did until the controller's
|
||||||
|
// next assertion moves it, and a permission that only allowed the new shape would refuse
|
||||||
|
// every node in the mesh for exactly as long as that took.
|
||||||
|
sub = []string{"mesh.node." + p.Node + ".declare",
|
||||||
|
"_DELIVER." + p.Node, "_DELIVER." + p.Node + ".>"}
|
||||||
|
|
||||||
case KindModule:
|
case KindModule:
|
||||||
// 1. Its own namespace: it publishes its events there and serves its tools there. Nothing
|
// 1. Its own namespace: it publishes its events there and serves its tools there. Nothing
|
||||||
@@ -323,7 +329,7 @@ func PermissionsFor(p Principal) (Permissions, error) {
|
|||||||
// take work over the new bus was refused the asking (2026-09-28).
|
// take work over the new bus was refused the asking (2026-09-28).
|
||||||
worker := "SEAT_" + upperSnake(s.Name) + "_worker"
|
worker := "SEAT_" + upperSnake(s.Name) + "_worker"
|
||||||
stream := seatStreamName(s.Name)
|
stream := seatStreamName(s.Name)
|
||||||
sub = append(sub, "_DELIVER."+worker)
|
sub = append(sub, "_DELIVER."+worker, "_DELIVER."+worker+".>")
|
||||||
pub = append(pub, "$JS.API.CONSUMER.INFO."+stream+"."+worker, "$JS.ACK."+stream+"."+worker+".>")
|
pub = append(pub, "$JS.API.CONSUMER.INFO."+stream+"."+worker, "$JS.ACK."+stream+"."+worker+".>")
|
||||||
for _, a := range s.Accepts {
|
for _, a := range s.Accepts {
|
||||||
sub = append(sub, seatSubject(s, "accept", a))
|
sub = append(sub, seatSubject(s, "accept", a))
|
||||||
|
|||||||
@@ -94,6 +94,24 @@ func MeshStreams() []Stream {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// DeliverSubjectFor is where a push consumer's messages land.
|
||||||
|
//
|
||||||
|
// **Per consumer, which means per stream as well as per name** (novox/hq 04-ISSUES/146). A push
|
||||||
|
// consumer delivers onto an ordinary subject, and everything subscribed to that subject gets a
|
||||||
|
// copy. The controller holds a consumer called `controller` on CONTROL and another called
|
||||||
|
// `controller` on EVENTS; while both were given `_DELIVER.controller`, the one process holding
|
||||||
|
// both subscriptions acted on every message twice — a joining machine was enrolled twice from one
|
||||||
|
// request, and the second enrolment minted a credential that replaced the one the machine had just
|
||||||
|
// been handed. Every report and every followed event doubled the same way, silently: nothing is
|
||||||
|
// redelivered, no count is wrong, the work simply happens twice.
|
||||||
|
//
|
||||||
|
// The stream belongs in it because the pair is what identifies a consumer — the server scopes a
|
||||||
|
// durable's name to its stream, and this subject was the one place that scoping was dropped. It
|
||||||
|
// stays inside what a controller may already subscribe (`_DELIVER.controller.>`).
|
||||||
|
func DeliverSubjectFor(c Consumer) string {
|
||||||
|
return "_DELIVER." + c.Name + "." + c.Stream
|
||||||
|
}
|
||||||
|
|
||||||
// An Asserter is the part of a JetStream connection stream assertion needs. Narrow on purpose: it
|
// An Asserter is the part of a JetStream connection stream assertion needs. Narrow on purpose: it
|
||||||
// keeps this testable without a server, and keeps the client library out of everything that only
|
// keeps this testable without a server, and keeps the client library out of everything that only
|
||||||
// wants to know what the streams are.
|
// wants to know what the streams are.
|
||||||
|
|||||||
@@ -233,3 +233,30 @@ func containsStep(steps []string, want string) bool {
|
|||||||
}
|
}
|
||||||
return false
|
return false
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// **Two consumers may share a name, and must not share a delivery subject** (novox/hq
|
||||||
|
// 04-ISSUES/146).
|
||||||
|
//
|
||||||
|
// A push consumer delivers onto an ordinary subject and everything subscribed to it gets a copy.
|
||||||
|
// The controller holds a consumer called `controller` on CONTROL and another called `controller` on
|
||||||
|
// EVENTS; while both were given `_DELIVER.controller`, the one process holding both subscriptions
|
||||||
|
// acted on every message twice — a joining machine enrolled twice from one request, with the second
|
||||||
|
// enrolment minting a credential that replaced the one the machine had just been handed.
|
||||||
|
//
|
||||||
|
// Checked here rather than against a server because it is a property of what the mesh asks for, and
|
||||||
|
// because the failure it produces is silent: every count is right, nothing is redelivered, and the
|
||||||
|
// work simply happens twice.
|
||||||
|
func TestNoTwoConsumersDeliverOntoTheSameSubject(t *testing.T) {
|
||||||
|
seen := map[string]string{}
|
||||||
|
for _, c := range MeshConsumers() {
|
||||||
|
if !c.Push && c.Queue == "" {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
subject := DeliverSubjectFor(c)
|
||||||
|
if other, taken := seen[subject]; taken {
|
||||||
|
t.Errorf("%s on %s and %s deliver onto %s, so whoever holds both acts on every "+
|
||||||
|
"message twice", c.Name, c.Stream, other, subject)
|
||||||
|
}
|
||||||
|
seen[subject] = c.Name + " on " + c.Stream
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
+2
-2
@@ -34,11 +34,11 @@ accounts {
|
|||||||
} }
|
} }
|
||||||
{ user: "node.one", password: "$2a$11$nnnnnnnnnnnnnnnnnnnnnn", permissions: {
|
{ user: "node.one", password: "$2a$11$nnnnnnnnnnnnnnnnnnnnnn", permissions: {
|
||||||
publish: { allow: ["$JS.ACK.NODES.one.>", "$JS.API.CONSUMER.INFO.NODES.one", "mesh.control.one.>"] }
|
publish: { allow: ["$JS.ACK.NODES.one.>", "$JS.API.CONSUMER.INFO.NODES.one", "mesh.control.one.>"] }
|
||||||
subscribe: { allow: ["_DELIVER.one", "_INBOX.node.one.>", "mesh.node.one.declare"] }
|
subscribe: { allow: ["_DELIVER.one", "_DELIVER.one.>", "_INBOX.node.one.>", "mesh.node.one.declare"] }
|
||||||
} }
|
} }
|
||||||
{ user: "one.telegram", password: "$2a$11$tttttttttttttttttttttt", permissions: {
|
{ user: "one.telegram", password: "$2a$11$tttttttttttttttttttttt", permissions: {
|
||||||
publish: { allow: ["$JS.ACK.EVENTS.one_telegram.>", "$JS.ACK.SEAT_TELEGRAM_SENDER.SEAT_TELEGRAM_SENDER_worker.>", "$JS.API.CONSUMER.INFO.EVENTS.one_telegram", "$JS.API.CONSUMER.INFO.SEAT_TELEGRAM_SENDER.SEAT_TELEGRAM_SENDER_worker", "$JS.API.CONSUMER.MSG.NEXT.EVENTS.one_telegram", "mesh.seat.telegram-sender.event.delivered", "mesh.seat.telegram-sender.event.failed"] }
|
publish: { allow: ["$JS.ACK.EVENTS.one_telegram.>", "$JS.ACK.SEAT_TELEGRAM_SENDER.SEAT_TELEGRAM_SENDER_worker.>", "$JS.API.CONSUMER.INFO.EVENTS.one_telegram", "$JS.API.CONSUMER.INFO.SEAT_TELEGRAM_SENDER.SEAT_TELEGRAM_SENDER_worker", "$JS.API.CONSUMER.MSG.NEXT.EVENTS.one_telegram", "mesh.seat.telegram-sender.event.delivered", "mesh.seat.telegram-sender.event.failed"] }
|
||||||
subscribe: { allow: ["_DELIVER.SEAT_TELEGRAM_SENDER_worker", "_INBOX.one.telegram.>", "mesh.mod.telegram.tool.>", "mesh.seat.telegram-sender.accept.send"] }
|
subscribe: { allow: ["_DELIVER.SEAT_TELEGRAM_SENDER_worker", "_DELIVER.SEAT_TELEGRAM_SENDER_worker.>", "_INBOX.one.telegram.>", "mesh.mod.telegram.tool.>", "mesh.seat.telegram-sender.accept.send"] }
|
||||||
allow_responses: { max: 1, ttl: "1m" }
|
allow_responses: { max: 1, ttl: "1m" }
|
||||||
} }
|
} }
|
||||||
{ user: "two.audit", password: "$2a$11$aaaaaaaaaaaaaaaaaaaaaa", permissions: {
|
{ user: "two.audit", password: "$2a$11$aaaaaaaaaaaaaaaaaaaaaa", permissions: {
|
||||||
|
|||||||
@@ -0,0 +1,71 @@
|
|||||||
|
package catalogue
|
||||||
|
|
||||||
|
import (
|
||||||
|
"os"
|
||||||
|
"strings"
|
||||||
|
"testing"
|
||||||
|
)
|
||||||
|
|
||||||
|
// **A machine trusts the mesh's authority because a module put its root there** (novox/hq ADR
|
||||||
|
// 0147, issue 129). The module carries a shell script and a unit, and both are worthless unless
|
||||||
|
// the mesh fills in where the authority is — which is the one thing about it the manifest cannot
|
||||||
|
// state, because the authority's address is a fact about the mesh and not about the module.
|
||||||
|
//
|
||||||
|
// So what is checked here is the rendering, not the parsing: the script the machine will run
|
||||||
|
// names the authority it was bound to, and the unit runs that script both ways. The verification
|
||||||
|
// itself — a plain client trusting an internal name on a machine holding this, and failing on one
|
||||||
|
// that does not — is the lab's, and cannot be had here.
|
||||||
|
func TestCaTrustRendersTheAuthorityItWasBoundTo(t *testing.T) {
|
||||||
|
raw, err := os.ReadFile("../../../mesh-catalog/modules/ca-trust/module.json")
|
||||||
|
if err != nil {
|
||||||
|
t.Skipf("the catalogue is not beside this checkout: %v", err)
|
||||||
|
}
|
||||||
|
m, err := ParseManifest(raw)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("the trust module does not parse:\n%v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
r := Resolution{
|
||||||
|
Node: "workstation",
|
||||||
|
Modules: []Manifest{m},
|
||||||
|
Needs: []Needed{{
|
||||||
|
Name: "internal-acme-ca", From: "anchor", At: "anchor.internal", For: "ca-trust",
|
||||||
|
Serves: map[string]any{
|
||||||
|
"port": float64(9000), "path": "/acme/acme/directory", "roots": "/roots.pem",
|
||||||
|
},
|
||||||
|
}},
|
||||||
|
}
|
||||||
|
out, err := r.Declaration(Rendering{})
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("the trust module could not be composed for a machine: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
script := fileNamed(out, "ca-trust.anchor")
|
||||||
|
if script == nil {
|
||||||
|
t.Fatalf("nothing writes the script the unit runs: %v", out)
|
||||||
|
}
|
||||||
|
body, _ := script["content"].(string)
|
||||||
|
if !strings.Contains(body, "https://anchor.internal:9000/roots.pem") {
|
||||||
|
t.Errorf("the script does not fetch from the authority it was bound to:\n%s", body)
|
||||||
|
}
|
||||||
|
if script["mode"] != "0755" {
|
||||||
|
t.Errorf("the script is written %v, which systemd cannot execute", script["mode"])
|
||||||
|
}
|
||||||
|
|
||||||
|
unit := fileNamed(out, "ca-trust.unit")
|
||||||
|
if unit == nil {
|
||||||
|
t.Fatalf("no unit: %v", out)
|
||||||
|
}
|
||||||
|
text, _ := unit["content"].(string)
|
||||||
|
// Both halves. A unit that only installs the anchor leaves a machine trusting an authority
|
||||||
|
// nobody assigned it to any more, which is the half issue 129 asked for by name.
|
||||||
|
for _, want := range []string{
|
||||||
|
"ExecStart=" + script["path"].(string) + " install",
|
||||||
|
"ExecStop=" + script["path"].(string) + " remove",
|
||||||
|
"RemainAfterExit=yes",
|
||||||
|
} {
|
||||||
|
if !strings.Contains(text, want) {
|
||||||
|
t.Errorf("the unit does not say %q:\n%s", want, text)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -1102,6 +1102,7 @@ func (r Resolution) contributions(settings SettingsBy, grants []Grant,
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, fmt.Errorf("%s contributing to %s: %w", m.Module, to, err)
|
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)
|
composeName(values, r.PublicDomain, r.At, reaches, endpointPorts(m), blocks)
|
||||||
out[to] = append(out[to], Contribution{From: m.Module, Values: values})
|
out[to] = append(out[to], Contribution{From: m.Module, Values: values})
|
||||||
}
|
}
|
||||||
@@ -1124,6 +1125,7 @@ func (r Resolution) contributions(settings SettingsBy, grants []Grant,
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, fmt.Errorf("%s contributing %s to %s: %w", m.Module, local, to, err)
|
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)
|
composeName(values, r.PublicDomain, r.At, reaches, endpointPorts(m), blocks)
|
||||||
out[to] = append(out[to], Contribution{From: m.Module, Values: values})
|
out[to] = append(out[to], Contribution{From: m.Module, Values: values})
|
||||||
}
|
}
|
||||||
@@ -1837,3 +1839,32 @@ func endpointPorts(m Manifest) map[string]int {
|
|||||||
}
|
}
|
||||||
return out
|
return out
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// portOfEndpoint fills in the port of the endpoint a contribution names, in place.
|
||||||
|
//
|
||||||
|
// **A contribution that names an endpoint must still carry that endpoint's port**, because everything
|
||||||
|
// downstream reads the port: the provider is told where to reach the consumer, and the machine-side
|
||||||
|
// redirection that turns a declared port into the number the machine published is keyed on it
|
||||||
|
// (atMachinePort). A route that named only its endpoint left the proxy with no port at all, and a
|
||||||
|
// proxy with no port has nothing to dial.
|
||||||
|
//
|
||||||
|
// Found before it shipped and after the catalogue had already been changed to name endpoints — the
|
||||||
|
// manifests were merged and the mesh had not yet picked them up, so nothing was broken yet. The
|
||||||
|
// declared port, not the machine one: the redirection happens later and is keyed on the declared
|
||||||
|
// number, so filling in the machine port here would be redirected a second time or not at all.
|
||||||
|
func portOfEndpoint(values map[string]any, ports map[string]int) {
|
||||||
|
if values == nil {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
if _, already := values["port"]; already {
|
||||||
|
// A route that says both is its own answer; the older shape repeated the port and is still read.
|
||||||
|
return
|
||||||
|
}
|
||||||
|
name, ok := values[RouteEndpoint].(string)
|
||||||
|
if !ok {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
if port, found := ports[strings.TrimSpace(name)]; found {
|
||||||
|
values["port"] = port
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -145,3 +145,61 @@ func TestAnUnnamedEndpointIsStillValid(t *testing.T) {
|
|||||||
t.Fatalf("a module with no route was refused: %v", got)
|
t.Fatalf("a module with no route was refused: %v", got)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// **A route that names an endpoint still carries that endpoint's port.**
|
||||||
|
//
|
||||||
|
// Everything downstream reads the port: the provider is told where to reach the consumer, and the
|
||||||
|
// redirection that turns a declared port into the number the machine published is keyed on it. A route
|
||||||
|
// naming only its endpoint left the proxy with no port, and a proxy with no port has nothing to dial.
|
||||||
|
//
|
||||||
|
// Caught after the catalogue had already been changed to name endpoints, and before the mesh picked
|
||||||
|
// those manifests up — which is the only reason nothing broke.
|
||||||
|
func TestARouteNamingAnEndpointStillCarriesItsPort(t *testing.T) {
|
||||||
|
m := aMediaServer()
|
||||||
|
r := Resolution{Node: "anchor", Modules: []Manifest{m},
|
||||||
|
PublicDomain: "example.test", At: "anchor.internal"}
|
||||||
|
given, err := r.contributions(nil, nil, nil)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
var saw bool
|
||||||
|
for _, c := range given["route"] {
|
||||||
|
saw = true
|
||||||
|
port, ok := asPort(c.Values["port"])
|
||||||
|
if !ok {
|
||||||
|
t.Fatalf("the route carries no port, so the proxy has nothing to dial: %v", c.Values)
|
||||||
|
}
|
||||||
|
if port != 80 {
|
||||||
|
t.Fatalf("the route carries port %d, want the web endpoint's 80", port)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if !saw {
|
||||||
|
t.Fatal("the module contributed no route")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// And the declared port, not the machine one: the redirection to where the machine published it
|
||||||
|
// happens later and is keyed on the declared number, so filling the machine port in here would be
|
||||||
|
// redirected twice or not at all.
|
||||||
|
func TestTheEndpointsDeclaredPortIsFilledInNotTheMachineOne(t *testing.T) {
|
||||||
|
m := aMediaServer()
|
||||||
|
values := map[string]any{RouteEndpoint: "web", "label": "media"}
|
||||||
|
portOfEndpoint(values, endpointPorts(m))
|
||||||
|
if got, _ := asPort(values["port"]); got != 80 {
|
||||||
|
t.Fatalf("filled in port %d, want the declared 80", got)
|
||||||
|
}
|
||||||
|
// Then the ordinary redirection puts it where the machine published it.
|
||||||
|
moved := atMachinePort(values, m.Module, map[string]map[int]int{"media": {80: 20009}})
|
||||||
|
if got, _ := asPort(moved["port"]); got != 20009 {
|
||||||
|
t.Fatalf("after redirection the port is %d, want the machine's 20009", got)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// A route that repeats a port keeps it, because that is the older shape and still read.
|
||||||
|
func TestARouteThatRepeatsItsPortKeepsIt(t *testing.T) {
|
||||||
|
values := map[string]any{RouteEndpoint: "web", "port": 8080}
|
||||||
|
portOfEndpoint(values, map[string]int{"web": 80})
|
||||||
|
if got, _ := asPort(values["port"]); got != 8080 {
|
||||||
|
t.Fatalf("the port it stated was overwritten with %d", got)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -265,6 +265,26 @@ func AsNftables(rules []Rule, mesh []string, outward bool, foundation []int,
|
|||||||
b.WriteString("\t\tct state established,related accept\n")
|
b.WriteString("\t\tct state established,related accept\n")
|
||||||
b.WriteString("\t\tct state invalid drop\n")
|
b.WriteString("\t\tct state invalid drop\n")
|
||||||
b.WriteString("\t\tiif lo accept\n")
|
b.WriteString("\t\tiif lo accept\n")
|
||||||
|
// **Anything on this machine may call anything on this machine.**
|
||||||
|
//
|
||||||
|
// Local is not a boundary this mesh draws. A service running here is callable by everything else
|
||||||
|
// running here, whatever form either takes — a package with a unit, a binary, a container. Whether
|
||||||
|
// a caller sits in a container was never meant to change the answer, and the only reason it did was
|
||||||
|
// that this chain asked about addresses: a caller on the machine carries the machine's address, a
|
||||||
|
// caller in one of its containers carries a bridge address, and a rule naming the former silently
|
||||||
|
// refused the latter.
|
||||||
|
//
|
||||||
|
// Measured: a module reaching its database on this machine's own name timed out for eleven hours
|
||||||
|
// while the machine itself could reach it, and the mesh called the machine healthy throughout
|
||||||
|
// (novox/hq 04-ISSUES/145).
|
||||||
|
//
|
||||||
|
// Asked by the link it arrives on rather than the address it comes from: anything that did not
|
||||||
|
// arrive from outside this machine, and did not arrive over the private network, is this machine's
|
||||||
|
// own. One rule for every service here, in place of a line per port that only ever covered the
|
||||||
|
// ports somebody remembered to think about.
|
||||||
|
if inward != "" {
|
||||||
|
b.WriteString(fmt.Sprintf("\t\tiifname != { %s } accept\n", inward))
|
||||||
|
}
|
||||||
b.WriteString("\t\ticmp type echo-request accept\n")
|
b.WriteString("\t\ticmp type echo-request accept\n")
|
||||||
b.WriteString("\t\ticmpv6 type { echo-request, nd-neighbor-solicit, nd-neighbor-advert, nd-router-advert } accept\n")
|
b.WriteString("\t\ticmpv6 type { echo-request, nd-neighbor-solicit, nd-neighbor-advert, nd-router-advert } accept\n")
|
||||||
|
|
||||||
|
|||||||
@@ -804,3 +804,86 @@ func TestSSHIsNeverLeftWithoutARule(t *testing.T) {
|
|||||||
t.Fatalf("a machine with no mesh addresses has no ssh rule, so adopting it locks it:\n%s", nft)
|
t.Fatalf("a machine with no mesh addresses has no ssh rule, so adopting it locks it:\n%s", nft)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// chainBody is one chain's own lines, so an assertion cannot be satisfied by an identical line in
|
||||||
|
// another chain.
|
||||||
|
//
|
||||||
|
// **Written because that happened.** The rule letting this machine's own callers through appears in the
|
||||||
|
// input chain and, in the same words, in the forward chain. A test asserting on the whole rendered file
|
||||||
|
// passed with the input chain's copy deleted — it was reading the forward chain's. ADR 0137's own tests
|
||||||
|
// say to assert per chain body for exactly this reason, and this file was not doing it.
|
||||||
|
func chainBody(t *testing.T, nft, chain string) string {
|
||||||
|
t.Helper()
|
||||||
|
open := "\tchain " + chain + " {"
|
||||||
|
i := strings.Index(nft, open)
|
||||||
|
if i < 0 {
|
||||||
|
t.Fatalf("no chain %q in:\n%s", chain, nft)
|
||||||
|
}
|
||||||
|
rest := nft[i+len(open):]
|
||||||
|
j := strings.Index(rest, "\n\t}")
|
||||||
|
if j < 0 {
|
||||||
|
t.Fatalf("chain %q does not close in:\n%s", chain, nft)
|
||||||
|
}
|
||||||
|
return rest[:j]
|
||||||
|
}
|
||||||
|
|
||||||
|
// **Anything on this machine may call anything on this machine.**
|
||||||
|
//
|
||||||
|
// Local is not a boundary this mesh draws, and whether a caller sits in a container was never meant to
|
||||||
|
// change the answer. It did, because the chain asked about addresses: a caller on the machine carries
|
||||||
|
// the machine's address and a caller in one of its containers carries a bridge address, so a rule
|
||||||
|
// naming the machines' own addresses silently refused every container on them.
|
||||||
|
//
|
||||||
|
// Measured: a module reaching its database on its own machine's name timed out for eleven hours while
|
||||||
|
// the machine itself could reach it (novox/hq 04-ISSUES/145).
|
||||||
|
func TestAnythingOnThisMachineMayCallAnythingOnIt(t *testing.T) {
|
||||||
|
nft := AsNftables(mustFilter(t, Resolution{Modules: []Manifest{
|
||||||
|
{Module: "store", Listens: []Listening{{Port: 5432, From: FromMesh}}},
|
||||||
|
{Module: "private", Listens: []Listening{{Port: 9999, From: FromMachine}}},
|
||||||
|
}}, nil), []string{"10.10.0.1", "10.10.0.2"}, false, nil, []string{"eth0"}, "mesh0")
|
||||||
|
|
||||||
|
// In the INPUT chain, which is where a call to a service on this machine arrives. The forward
|
||||||
|
// chain carries the same line in the same words, so asserting on the whole file proves nothing.
|
||||||
|
input := chainBody(t, nft, "input")
|
||||||
|
if !strings.Contains(input, `iifname != { "eth0", "mesh0" } accept`) {
|
||||||
|
t.Fatalf("a caller on this machine cannot reach a service on it:\n%s", input)
|
||||||
|
}
|
||||||
|
// One rule, for every service here — not a line per port that only covers the ports somebody
|
||||||
|
// remembered to think about.
|
||||||
|
if strings.Contains(input, `iifname != { "eth0", "mesh0" } tcp dport 5432`) {
|
||||||
|
t.Fatalf("the local allowance is still written per port:\n%s", input)
|
||||||
|
}
|
||||||
|
// And the private network still reaches what is exposed to it, which is a different question.
|
||||||
|
if !strings.Contains(input, "ip saddr { 10.10.0.1, 10.10.0.2 } tcp dport 5432 accept") {
|
||||||
|
t.Fatalf("the private network no longer reaches a service exposed to it:\n%s", input)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// The three reaches, as three lines. This is the whole of what the filter says about who may call what.
|
||||||
|
func TestTheThreeReachesAreThreeLines(t *testing.T) {
|
||||||
|
nft := AsNftables(mustFilter(t, Resolution{Modules: []Manifest{
|
||||||
|
{Module: "internal-only", Listens: []Listening{{Port: 5432, From: FromMesh}}},
|
||||||
|
{Module: "public", Listens: []Listening{{Port: 443, From: FromEverywhere}}},
|
||||||
|
}}, nil), []string{"10.10.0.1"}, false, nil, []string{"eth0"}, "mesh0")
|
||||||
|
|
||||||
|
input := chainBody(t, nft, "input")
|
||||||
|
for what, want := range map[string]string{
|
||||||
|
"on this machine": `iifname != { "eth0", "mesh0" } accept`,
|
||||||
|
"over the private network": "ip saddr { 10.10.0.1 } tcp dport 5432 accept",
|
||||||
|
"from anywhere": "tcp dport 443 accept",
|
||||||
|
} {
|
||||||
|
if !strings.Contains(input, want) {
|
||||||
|
t.Fatalf("a caller %s cannot reach what is exposed to it (%q):\n%s", what, want, input)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// A port open to everything needs no such line — it is already open to a guest.
|
||||||
|
func TestAPublicPortNeedsNoGuestLine(t *testing.T) {
|
||||||
|
nft := AsNftables(mustFilter(t, Resolution{Modules: []Manifest{
|
||||||
|
{Module: "web", Listens: []Listening{{Port: 443, From: FromEverywhere}}},
|
||||||
|
}}, nil), []string{"10.10.0.1"}, false, nil, []string{"eth0"}, "mesh0")
|
||||||
|
if strings.Count(nft, `iifname != { "eth0", "mesh0" } tcp dport 443`) != 0 {
|
||||||
|
t.Fatalf("a public port was given a guest line it does not need:\n%s", nft)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -8,6 +8,7 @@ import (
|
|||||||
|
|
||||||
"github.com/novox/mesh-controller/internal/broker"
|
"github.com/novox/mesh-controller/internal/broker"
|
||||||
"github.com/novox/mesh-controller/internal/catalogue"
|
"github.com/novox/mesh-controller/internal/catalogue"
|
||||||
|
"golang.org/x/crypto/bcrypt"
|
||||||
)
|
)
|
||||||
|
|
||||||
// Reading the bus's user list out of the mesh's records, against a real store.
|
// Reading the bus's user list out of the mesh's records, against a real store.
|
||||||
@@ -164,3 +165,63 @@ func granted(all []string, one string) bool {
|
|||||||
}
|
}
|
||||||
return false
|
return false
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// **A token is an account on the bus, or it is a string nothing accepts** (novox/hq 04-ISSUES/146).
|
||||||
|
//
|
||||||
|
// The composed list names an enrolment user for every machine with a live token, and nothing minted
|
||||||
|
// a credential for it — so the composer left it out as a user with no password, and every enrolment
|
||||||
|
// since the mesh moved to this bus was refused by the server before the mesh heard of it. Nothing
|
||||||
|
// caught it because nothing had enrolled since.
|
||||||
|
//
|
||||||
|
// The password cannot be minted, because it is the token's own secret: the machine will present
|
||||||
|
// exactly that string. So this checks the two halves that make the account usable — that a row
|
||||||
|
// exists under the name the composer asks for, and that the secret handed out is what that row
|
||||||
|
// accepts.
|
||||||
|
func TestIssuingATokenRecordsTheAccountItIsThePasswordOf(t *testing.T) {
|
||||||
|
inv, ctx := aMeshWith(t)
|
||||||
|
if _, err := inv.AddNode(ctx, "joiner"); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
issued, err := inv.IssueToken(ctx, "joiner", time.Hour)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
|
||||||
|
name := broker.Principal{Kind: broker.KindEnrolment, Node: "joiner"}.Username()
|
||||||
|
users, err := inv.BusUsers(ctx)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
user, has := users[name]
|
||||||
|
if !has {
|
||||||
|
t.Fatalf("no bus account for %q; the composer would leave the enrolment out and the "+
|
||||||
|
"machine would be refused before the mesh heard of it: %v", name, users)
|
||||||
|
}
|
||||||
|
if user.Kind != BusEnrolment || user.Node != "joiner" {
|
||||||
|
t.Errorf("the account is %+v, not this node's enrolment", user)
|
||||||
|
}
|
||||||
|
if err := bcrypt.CompareHashAndPassword([]byte(user.PasswordHash), []byte(issued.Secret)); err != nil {
|
||||||
|
t.Error("the account does not accept the secret the token carries, so presenting the " +
|
||||||
|
"token would be refused by the server")
|
||||||
|
}
|
||||||
|
|
||||||
|
// And the composition contains it, which is the thing the server reads.
|
||||||
|
records, err := inv.BusRecords(ctx)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
derived, err := broker.Users(records)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
hashes := map[string]string{}
|
||||||
|
for n, u := range users {
|
||||||
|
hashes[n] = u.PasswordHash
|
||||||
|
}
|
||||||
|
_, missing := broker.WithPasswords(derived, hashes)
|
||||||
|
for _, m := range missing {
|
||||||
|
if m == name {
|
||||||
|
t.Fatal("the enrolment user is composed without a password, which is a user nobody can be")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -51,6 +51,40 @@ const (
|
|||||||
// reply, into a module's sealed environment — and the mesh keeps only the hash, so a credential is
|
// reply, into a module's sealed environment — and the mesh keeps only the hash, so a credential is
|
||||||
// never recoverable from the store. A caller that loses it must mint again, which is a rotation and
|
// never recoverable from the store. A caller that loses it must mint again, which is a rotation and
|
||||||
// is meant to feel like one.
|
// is meant to feel like one.
|
||||||
|
// RecordBusPassword records a hash for a password the caller already holds.
|
||||||
|
//
|
||||||
|
// **For the one credential the mesh does not choose**: an enrolment token's secret is the password
|
||||||
|
// of the user that presents it (novox/hq ADR 0004, design 25 §6), so the token cannot be given a
|
||||||
|
// minted password — it already has one, and the machine will connect with exactly that string.
|
||||||
|
// Everything else goes through Mint, which chooses and returns the plaintext once.
|
||||||
|
func (i *Inventory) RecordBusPassword(ctx context.Context, u BusUser, password string) error {
|
||||||
|
if u.Username == "" || u.Kind == "" {
|
||||||
|
return errors.New("a bus user needs a username and a kind")
|
||||||
|
}
|
||||||
|
if password == "" {
|
||||||
|
return errors.New("a bus user needs a password")
|
||||||
|
}
|
||||||
|
hash, err := bcrypt.GenerateFromPassword([]byte(password), bcrypt.DefaultCost)
|
||||||
|
if err != nil {
|
||||||
|
return fmt.Errorf("cannot hash a bus password: %w", err)
|
||||||
|
}
|
||||||
|
return i.writeBusUser(ctx, u, string(hash))
|
||||||
|
}
|
||||||
|
|
||||||
|
// writeBusUser is the row, whoever chose the password.
|
||||||
|
func (i *Inventory) writeBusUser(ctx context.Context, u BusUser, hash string) error {
|
||||||
|
if _, err := i.store.Pool().Exec(ctx,
|
||||||
|
`insert into bus_user (username, kind, node, module, password_hash)
|
||||||
|
values ($1, $2, $3, $4, $5)
|
||||||
|
on conflict (username) do update
|
||||||
|
set kind = excluded.kind, node = excluded.node, module = excluded.module,
|
||||||
|
password_hash = excluded.password_hash, minted_at = now()`,
|
||||||
|
u.Username, u.Kind, u.Node, u.Module, hash); err != nil {
|
||||||
|
return fmt.Errorf("cannot record the bus user %s: %w", u.Username, err)
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
func (i *Inventory) MintBusPassword(ctx context.Context, u BusUser) (string, error) {
|
func (i *Inventory) MintBusPassword(ctx context.Context, u BusUser) (string, error) {
|
||||||
if u.Username == "" || u.Kind == "" {
|
if u.Username == "" || u.Kind == "" {
|
||||||
return "", errors.New("a bus user needs a username and a kind")
|
return "", errors.New("a bus user needs a username and a kind")
|
||||||
@@ -69,14 +103,8 @@ func (i *Inventory) MintBusPassword(ctx context.Context, u BusUser) (string, err
|
|||||||
return "", fmt.Errorf("cannot hash a bus password: %w", err)
|
return "", fmt.Errorf("cannot hash a bus password: %w", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
if _, err := i.store.Pool().Exec(ctx,
|
if err := i.writeBusUser(ctx, u, string(hash)); err != nil {
|
||||||
`insert into bus_user (username, kind, node, module, password_hash)
|
return "", err
|
||||||
values ($1, $2, $3, $4, $5)
|
|
||||||
on conflict (username) do update
|
|
||||||
set kind = excluded.kind, node = excluded.node, module = excluded.module,
|
|
||||||
password_hash = excluded.password_hash, minted_at = now()`,
|
|
||||||
u.Username, u.Kind, u.Node, u.Module, string(hash)); err != nil {
|
|
||||||
return "", fmt.Errorf("cannot record the bus user %s: %w", u.Username, err)
|
|
||||||
}
|
}
|
||||||
return password, nil
|
return password, nil
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -14,6 +14,7 @@ import (
|
|||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/jackc/pgx/v5"
|
"github.com/jackc/pgx/v5"
|
||||||
|
"github.com/novox/mesh-controller/internal/broker"
|
||||||
"github.com/novox/mesh-controller/internal/store"
|
"github.com/novox/mesh-controller/internal/store"
|
||||||
)
|
)
|
||||||
|
|
||||||
@@ -263,6 +264,29 @@ func (i *Inventory) IssueToken(ctx context.Context, nodeName string, validFor ti
|
|||||||
return Issued{}, err
|
return Issued{}, err
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// **And the account that secret is the password of** (novox/hq 04-ISSUES/146). The composed
|
||||||
|
// user list names an enrolment user for every node with a live token, and nothing minted a
|
||||||
|
// credential for it — so the composer left it out as a user with no password and every
|
||||||
|
// enrolment was refused by the server before the mesh heard of it.
|
||||||
|
//
|
||||||
|
// Recorded rather than minted: the token's secret IS the password, which is what lets a
|
||||||
|
// machine's first connection be authenticated by the thing it is enrolling with. It cannot be
|
||||||
|
// chosen here, because it has already been handed to whoever will present it.
|
||||||
|
//
|
||||||
|
// Outside the transaction on purpose. The token is what the mesh promised; a credential that
|
||||||
|
// the next composition rewrites anyway is not worth failing an issue over, and a token with no
|
||||||
|
// account is recoverable by issuing another, while an account with no token is a user nobody
|
||||||
|
// can be.
|
||||||
|
if err := i.RecordBusPassword(ctx, BusUser{
|
||||||
|
Username: broker.Principal{Kind: broker.KindEnrolment, Node: node.Name}.Username(),
|
||||||
|
Kind: BusEnrolment,
|
||||||
|
Node: node.Name,
|
||||||
|
}, secret); err != nil {
|
||||||
|
return Issued{}, fmt.Errorf(
|
||||||
|
"the token for %s was issued and the bus account it is the password of was not "+
|
||||||
|
"recorded, so this token cannot connect: %w", node.Name, err)
|
||||||
|
}
|
||||||
|
|
||||||
return Issued{Node: node, Secret: secret, Expires: expires}, nil
|
return Issued{Node: node, Secret: secret, Expires: expires}, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user