Compare commits
1
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
30362118a1 |
@@ -17,7 +17,12 @@ package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"crypto/sha256"
|
||||
"crypto/tls"
|
||||
"crypto/x509"
|
||||
"encoding/hex"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"net/url"
|
||||
"os"
|
||||
@@ -25,6 +30,8 @@ import (
|
||||
"strings"
|
||||
"syscall"
|
||||
|
||||
amqp "github.com/rabbitmq/amqp091-go"
|
||||
|
||||
"github.com/novox/mesh-controller/internal/broker"
|
||||
"github.com/novox/mesh-controller/internal/builder"
|
||||
"github.com/novox/mesh-controller/internal/link"
|
||||
@@ -45,6 +52,7 @@ const usage = `mesh-builder — builds modules for the mesh
|
||||
It consumes build requests and answers with what it made. Nothing is listened on and nothing
|
||||
is dialled except the broker.
|
||||
|
||||
MESH_BROKER_AMQP where the broker is, with this builder's own credential
|
||||
MESH_BROKER_FILE a file the mesh sealed to this machine holding the same
|
||||
MESH_REGISTRY host:port to publish artifacts to, when the mesh has not said
|
||||
MESH_BINDING a file the mesh wrote saying where the artifact store is
|
||||
@@ -127,17 +135,42 @@ func run() error {
|
||||
// machine told about both would take work from one and answer on the other, and every log line would
|
||||
// say it was fine.
|
||||
func takeWorkFrom(credential Credential, on string) (link.BuildMachine, error) {
|
||||
// **The credential names the bus, and there is one** (novox/hq ADR 0131, design 28 task 5.5).
|
||||
// A credential for the mesh's bus carries user, password and fingerprint beside the address,
|
||||
// and that is enough to dial it, pinned.
|
||||
if !credential.onTheNewBus() {
|
||||
return nil, fmt.Errorf("the credential at hand names %q, which is not the mesh's bus", credential.URL)
|
||||
// **The credential decides, before any variable does.** A machine moved to the new bus was
|
||||
// handed a credential for it and nothing else changed in its environment; that credential
|
||||
// names the bus by scheme, so it is enough to know which bus to take work from.
|
||||
if credential.onTheNewBus() {
|
||||
js, err := broker.DialPinned(credential.natsURL(), credential.Fingerprint)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return link.MachineOverNATS(js, on), nil
|
||||
}
|
||||
js, err := broker.DialPinned(credential.natsURL(), credential.Fingerprint)
|
||||
address, onNATS, err := broker.OnNATS()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return link.MachineOverNATS(js, on), nil
|
||||
if err := broker.MustBeOneBus(credential.URL, address); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if onNATS {
|
||||
js, err := broker.Dial(address)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("cannot reach the bus at %s: %w", address, err)
|
||||
}
|
||||
return link.MachineOverNATS(js, on), nil
|
||||
}
|
||||
|
||||
conn, err := dial(credential)
|
||||
if err != nil {
|
||||
// Not quoted back: the URL carries this builder's broker password.
|
||||
return nil, fmt.Errorf("cannot reach the broker: %w", err)
|
||||
}
|
||||
channel, err := conn.Channel()
|
||||
if err != nil {
|
||||
conn.Close()
|
||||
return nil, err
|
||||
}
|
||||
return link.MachineOverCurrent(conn, channel, on), nil
|
||||
}
|
||||
|
||||
// answer does one build and says what happened, whichever way it went.
|
||||
@@ -419,8 +452,13 @@ func brokerFrom() (Credential, error) {
|
||||
// broker is then verified against whatever this machine already trusts.
|
||||
return Credential{URL: said}, nil
|
||||
}
|
||||
return Credential{}, fmt.Errorf(
|
||||
"no MESH_BROKER_FILE: a build machine with no credential for the bus has nothing to build")
|
||||
url := strings.TrimSpace(os.Getenv("MESH_BROKER_AMQP"))
|
||||
if url == "" {
|
||||
return Credential{}, fmt.Errorf(
|
||||
"neither MESH_BROKER_FILE nor MESH_BROKER_AMQP: a builder with no broker has " +
|
||||
"nothing to build")
|
||||
}
|
||||
return Credential{URL: url}, nil
|
||||
}
|
||||
|
||||
// Credential is what a build machine is given so it can reach the broker.
|
||||
@@ -453,3 +491,43 @@ func (c Credential) natsURL() string {
|
||||
}
|
||||
return "nats://" + c.User + ":" + c.Password + "@" + rest
|
||||
}
|
||||
|
||||
// dial opens the connection, pinning the broker's certificate when there is one to pin.
|
||||
func dial(held Credential) (*amqp.Connection, error) {
|
||||
if held.Fingerprint == "" {
|
||||
return amqp.Dial(held.URL)
|
||||
}
|
||||
return amqp.DialTLS(held.URL, pinning(held.Fingerprint))
|
||||
}
|
||||
|
||||
// pinning is a TLS configuration that trusts exactly one certificate.
|
||||
//
|
||||
// InsecureSkipVerify with a VerifyPeerCertificate is **pinning, not skipping**: the standard chain
|
||||
// check is replaced, not removed, and what replaces it is stricter — one certificate is accepted
|
||||
// rather than every certificate a public authority would sign.
|
||||
//
|
||||
// Its own function so a test can drive it against a real handshake. A pin check that is only ever
|
||||
// exercised through a broker is a pin check nothing tests.
|
||||
func pinning(fingerprint string) *tls.Config {
|
||||
return &tls.Config{
|
||||
InsecureSkipVerify: true,
|
||||
VerifyPeerCertificate: func(raw [][]byte, _ [][]*x509.Certificate) error {
|
||||
if len(raw) == 0 {
|
||||
return errors.New("the broker presented no certificate")
|
||||
}
|
||||
// The leaf, and in the same spelling the mesh writes it — `sha256:` and 64 hex
|
||||
// characters. Comparing a bare digest against a written fingerprint never matches,
|
||||
// and the failure is indistinguishable from being pointed at the wrong broker.
|
||||
sum := sha256.Sum256(raw[0])
|
||||
got := "sha256:" + hex.EncodeToString(sum[:])
|
||||
if got != fingerprint {
|
||||
return fmt.Errorf(
|
||||
"this is not the broker this builder was told about\n expected %s\n "+
|
||||
"got %s\nEither this mesh's broker was replaced, or this builder is "+
|
||||
"being pointed at something else. Retrying will not help",
|
||||
fingerprint, got)
|
||||
}
|
||||
return nil
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
@@ -15,8 +15,6 @@ import (
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/novox/mesh-controller/internal/broker"
|
||||
)
|
||||
|
||||
// Where a builder publishes.
|
||||
@@ -195,9 +193,9 @@ func TestThePinIsComparedInTheSpellingTheMeshWritesIt(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
// handshakeWith runs the pin check the builder dials with against an address.
|
||||
// handshakeWith runs the builder's own pin check against an address.
|
||||
func handshakeWith(address, pin string) error {
|
||||
conn, err := tls.Dial("tcp", address, broker.PinnedToFingerprint(pin))
|
||||
conn, err := tls.Dial("tcp", address, pinning(pin))
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
@@ -43,7 +43,7 @@ func askCommand(ctx context.Context, args []string) error {
|
||||
}
|
||||
defer server.Close()
|
||||
|
||||
answer, err := link.Ask(ctx, server.Bus(), module, tool, arguments, *wait)
|
||||
answer, err := link.Ask(ctx, server.Channel(), module, tool, arguments, *wait)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
@@ -2,6 +2,8 @@ package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"crypto/rand"
|
||||
"encoding/base64"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"flag"
|
||||
@@ -48,7 +50,7 @@ func buildOn(ctx context.Context, base string, wait time.Duration) error {
|
||||
continue
|
||||
}
|
||||
for _, b := range e.Manifest.Build.On {
|
||||
if standsOnModule(b, base) {
|
||||
if b.Module == base {
|
||||
on = append(on, e)
|
||||
break
|
||||
}
|
||||
@@ -259,45 +261,85 @@ func builderCommand(ctx context.Context, args []string) error {
|
||||
}
|
||||
name := positionals[1]
|
||||
|
||||
// **The build machine's credential is a module's credential** (novox/hq ADR 0131, design 28
|
||||
// task 5.5): minted into the mesh's records and sealed to the machine as the builder module's
|
||||
// broker secret, usable at the next push — the same act `module issue` performs, and the same
|
||||
// account the composed user list carries. Nothing is created on a server; the bus reads the list.
|
||||
open, err := openStores(ctx)
|
||||
management, err := broker.ManagementFromEnvironment()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer open.Close()
|
||||
inv := open.inventory
|
||||
shelf, err := inv.Catalogue(ctx)
|
||||
if err != nil {
|
||||
|
||||
// The same shape of secret a token carries: enough entropy that guessing is not a strategy,
|
||||
// and safe to put in a URL because that is where it goes.
|
||||
raw := make([]byte, 32)
|
||||
if _, err := rand.Read(raw); err != nil {
|
||||
return err
|
||||
}
|
||||
m, known := shelf[*module]
|
||||
if !known {
|
||||
return fmt.Errorf("%s is not in the catalogue; `module add` it first", *module)
|
||||
password := base64.RawURLEncoding.EncodeToString(raw)
|
||||
if err := management.CreateBuilderAccount(ctx, name, password); err != nil {
|
||||
return err
|
||||
}
|
||||
fmt.Printf("build machine %s: ", name)
|
||||
node := *forNode
|
||||
if node == "" {
|
||||
entries, err := inv.Catalogued(ctx)
|
||||
|
||||
fmt.Printf("broker account %s created, scoped to the %s queue and the %s exchange\n\n",
|
||||
name, link.BuildQueue, link.Exchange)
|
||||
|
||||
if *forNode != "" {
|
||||
known, err := broker.FromEnvironment()
|
||||
if err != nil {
|
||||
return fmt.Errorf("cannot deliver a credential without knowing where the broker is: %w", err)
|
||||
}
|
||||
inv, err := openInventory(ctx)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
for _, e := range entries {
|
||||
if e.Manifest.Module == *module && len(e.On) > 0 {
|
||||
node = e.On[0]
|
||||
}
|
||||
defer inv.Close()
|
||||
|
||||
brokerAddr, err := brokerReachableAt(ctx, inv, known, *forNode)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
// The URL and what verifies the broker, together. A mesh's broker presents a certificate
|
||||
// of the mesh's own, which is in no public trust store — so a URL on its own reaches only
|
||||
// a broker somebody else vouches for, and the connection fails at TLS with an error about
|
||||
// an unknown authority rather than about a missing pin.
|
||||
//
|
||||
// **The same two facts a node's token carries** (novox/hq ADR 0004), delivered the same
|
||||
// way: out of band relative to the broker, so what is trusted does not come from the thing
|
||||
// being trusted.
|
||||
held, err := json.Marshal(struct {
|
||||
URL string `json:"url"`
|
||||
Fingerprint string `json:"fingerprint,omitempty"`
|
||||
}{
|
||||
URL: fmt.Sprintf("amqps://%s:%s@%s/", name, password, brokerAddr),
|
||||
Fingerprint: known.Fingerprint,
|
||||
})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if err := inv.AcceptSecretForModule(ctx, *forNode, *module, "broker", string(held)); err != nil {
|
||||
return err
|
||||
}
|
||||
// Not printed. It is sealed to that machine and the mesh cannot read it back, which is
|
||||
// the whole point — printing it here would put the one copy that matters on a terminal.
|
||||
fmt.Printf(" sealed to %s, for the %s module. It arrives with the next push.\n",
|
||||
*forNode, *module)
|
||||
fmt.Printf(" run `push %s` to send it\n", *forNode)
|
||||
return nil
|
||||
}
|
||||
if node == "" {
|
||||
return fmt.Errorf("%s is assigned nowhere; `assign <machine> %s` first, or say --node", *module, *module)
|
||||
|
||||
// The whole line only when the address is known. A URL with a placeholder where the host
|
||||
// should be is a URL somebody pastes and then debugs, and the placeholder is the last thing
|
||||
// they look at.
|
||||
if known, err := broker.FromEnvironment(); err == nil {
|
||||
fmt.Printf(" MESH_BROKER_AMQP=amqps://%s:%s@%s/\n\n", name, password, known.Address)
|
||||
} else {
|
||||
fmt.Printf(" the password is %s\n\n", password)
|
||||
fmt.Printf(" This control plane has no %s, so it cannot say where the broker is.\n"+
|
||||
" Put the password in MESH_BROKER_AMQP on the build machine.\n\n",
|
||||
broker.AddressVar)
|
||||
}
|
||||
address, err := broker.BusAddress()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
return issueOnTheNewBus(ctx, inv, m, node, address)
|
||||
// Shown once, like a token, and for the same reason: what is stored is the broker's own hash
|
||||
// of it, and a control plane that could show it back would be a control plane that holds it.
|
||||
fmt.Println("This is the only time it is shown.")
|
||||
return nil
|
||||
}
|
||||
|
||||
// buildBehind builds every module the mesh holds older than its source has.
|
||||
@@ -592,10 +634,16 @@ func heldBy(ctx context.Context) map[string]string {
|
||||
// **One place chooses**, as everywhere else the bus change went (novox/hq ADR 0116 step 5). On the bus
|
||||
// the mesh runs on today this needs the controller's own connection, so it is handed one; on the bus
|
||||
// being built it dials, because a build request is a one-shot and holds nothing else.
|
||||
func askOver(_ *link.Server) (link.Builders, error) {
|
||||
address, err := broker.BusAddress()
|
||||
func askOver(server *link.Server) (link.Builders, error) {
|
||||
address, onNATS, err := broker.OnNATS()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return link.BuildsOverNATS(address)
|
||||
if err := broker.MustBeOneBus(os.Getenv(broker.AMQPVarName), address); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if onNATS {
|
||||
return link.BuildsOverNATS(address)
|
||||
}
|
||||
return link.BuildsOverCurrent(server.Channel()), nil
|
||||
}
|
||||
|
||||
@@ -2,6 +2,8 @@ package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"crypto/rand"
|
||||
"encoding/base64"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"flag"
|
||||
@@ -259,11 +261,77 @@ func moduleCommand(ctx context.Context, args []string) error {
|
||||
// made in entirely different ways: on the bus the mesh runs on today an account is a
|
||||
// management call, and on the bus being built it is a row the next composition writes into
|
||||
// the server's user list (novox/hq design 25 §4).
|
||||
busAddress, err := broker.BusAddress()
|
||||
busAddress, onNATS, err := broker.OnNATS()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
return issueOnTheNewBus(ctx, inv, m, *forNode, busAddress)
|
||||
if err := broker.MustBeOneBus(os.Getenv(broker.AMQPVarName), busAddress); err != nil {
|
||||
return err
|
||||
}
|
||||
if onNATS {
|
||||
return issueOnTheNewBus(ctx, inv, m, *forNode, busAddress)
|
||||
}
|
||||
|
||||
management, err := broker.ManagementFromEnvironment()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
// The foundation owns the bus; make sure it exists before a module binds onto it.
|
||||
if err := management.EnsureEventExchanges(ctx); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
secret := make([]byte, 32)
|
||||
if _, err := rand.Read(secret); err != nil {
|
||||
return err
|
||||
}
|
||||
password := base64.RawURLEncoding.EncodeToString(secret)
|
||||
account, err := management.CreateModuleAccount(ctx, *forNode, module, password, m.Emits, m.Consumes)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
// A consumer's queue, with its dead-letter, is the foundation's to declare — its own account
|
||||
// may not (ADR 0043). Made now, so it exists before the module binds onto it.
|
||||
if len(m.Consumes) > 0 {
|
||||
if err := management.EnsureModuleQueue(ctx, *forNode, module); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
||||
known, err := broker.FromEnvironment()
|
||||
if err != nil {
|
||||
return fmt.Errorf("cannot deliver a credential without knowing where the broker is: %w", err)
|
||||
}
|
||||
brokerAddr, err := brokerReachableAt(ctx, inv, known, *forNode)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
// The URL and what verifies the broker, together — a mesh's broker presents its own
|
||||
// certificate, in no public trust store, so a URL alone fails at TLS (as `builder issue`).
|
||||
held, err := json.Marshal(struct {
|
||||
URL string `json:"url"`
|
||||
Fingerprint string `json:"fingerprint,omitempty"`
|
||||
Node string `json:"node"`
|
||||
Module string `json:"module"`
|
||||
}{
|
||||
URL: fmt.Sprintf("amqps://%s:%s@%s/", account, password, brokerAddr),
|
||||
Fingerprint: known.Fingerprint,
|
||||
// The node and module the account is for, so the runtime names its queue as the mesh
|
||||
// scoped it (<node>.<module>.events) without a manifest having to interpolate a node.
|
||||
Node: *forNode,
|
||||
Module: module,
|
||||
})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if err := inv.AcceptSecretForModule(ctx, *forNode, module, "broker", string(held)); err != nil {
|
||||
return err
|
||||
}
|
||||
fmt.Printf("broker account %s created for %s, scoped to what it emits and consumes\n",
|
||||
account, module)
|
||||
fmt.Printf(" sealed to %s. It arrives with the next push — `push %s` to send it\n",
|
||||
*forNode, *forNode)
|
||||
return nil
|
||||
|
||||
default:
|
||||
return fmt.Errorf("module has no %q; it has add, list, moved, forget and issue", args[0])
|
||||
|
||||
@@ -10,6 +10,7 @@ import (
|
||||
|
||||
"github.com/novox/mesh-controller/internal/broker"
|
||||
"github.com/novox/mesh-controller/internal/inventory"
|
||||
"github.com/novox/mesh-controller/internal/link"
|
||||
"github.com/novox/mesh-controller/internal/token"
|
||||
)
|
||||
|
||||
@@ -257,6 +258,15 @@ func tokenCommand(ctx context.Context, args []string) error {
|
||||
// chicken-and-egg entirely: the mesh runs the broker, so a joining node's credentials can
|
||||
// exist before it does. The one-time secret IS the password, so a node's first connection is
|
||||
// already authenticated and enrolment is what happens over it.
|
||||
if management, err := broker.ManagementFromEnvironment(); err == nil {
|
||||
if err := management.CreateNodeAccount(ctx, issued.Node.Name, issued.Secret); err != nil {
|
||||
return err
|
||||
}
|
||||
fmt.Printf("broker account %s created, scoped to %s and the %s exchange\n\n",
|
||||
issued.Node.Name, link.QueueFor(issued.Node.Name), link.Exchange)
|
||||
} else if !errors.Is(err, broker.ErrNotConfigured) {
|
||||
return err
|
||||
}
|
||||
|
||||
made := token.Token{Node: issued.Node.Name, Signer: key.Public, Secret: issued.Secret,
|
||||
Adopted: issued.Node.Adopted}
|
||||
|
||||
@@ -42,10 +42,16 @@ func reportUnhostable(node string, plan catalogue.Resolution) {
|
||||
// variable moves it). The streams and this controller's consumers are raised first on the new bus,
|
||||
// so nothing served here finds them missing.
|
||||
func connectLink(ctx context.Context, inv *inventory.Inventory, enroller link.Enroller, listener link.Listener) (*link.Server, error) {
|
||||
busAddress, err := broker.BusAddress()
|
||||
busAddress, onNATS, err := broker.OnNATS()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if err := broker.MustBeOneBus(os.Getenv(broker.AMQPVarName), busAddress); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if !onNATS {
|
||||
return link.Connect(enroller, listener)
|
||||
}
|
||||
if inv != nil {
|
||||
if err := raiseTheBus(ctx, inv, busAddress); err != nil {
|
||||
return nil, err
|
||||
@@ -82,6 +88,11 @@ func serve(ctx context.Context) error {
|
||||
}
|
||||
fmt.Printf("signing as %s\n", key.Fingerprint()[:16])
|
||||
|
||||
management, err := broker.ManagementFromEnvironment()
|
||||
if err != nil && !errors.Is(err, broker.ErrNotConfigured) {
|
||||
return err
|
||||
}
|
||||
|
||||
// Where the broker is and what to expect there, so a node can be told how to come back
|
||||
// without a person and a new token.
|
||||
known, err := broker.FromEnvironment()
|
||||
@@ -96,9 +107,16 @@ func serve(ctx context.Context) error {
|
||||
// **Which bus this mesh is on, read once** (novox/hq ADR 0116 step 5). Both clients ship; both
|
||||
// being live is refused, because a mesh half on each is one where a declaration goes out on one
|
||||
// and the report comes back on the other, and every component logs success while it happens.
|
||||
busAddress, onNATS, err := broker.OnNATS()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if err := broker.MustBeOneBus(os.Getenv(broker.AMQPVarName), busAddress); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
work := link.Enrolment{Inventory: inv, Identity: ident, Broker: known,
|
||||
OnNATS: true}
|
||||
work := link.Enrolment{Inventory: inv, Identity: ident, Management: management, Broker: known,
|
||||
OnNATS: onNATS}
|
||||
server, err := connectLink(ctx, inv, work, work)
|
||||
if err != nil {
|
||||
return err
|
||||
|
||||
@@ -8,7 +8,6 @@ import (
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/novox/mesh-controller/internal/catalogue"
|
||||
"github.com/novox/mesh-controller/internal/inventory"
|
||||
"github.com/novox/mesh-controller/internal/link"
|
||||
)
|
||||
@@ -311,8 +310,11 @@ func orderByBases(entries []inventory.Entry) []inventory.Entry {
|
||||
seen[name] = true
|
||||
if e.Manifest.Build != nil {
|
||||
for _, on := range e.Manifest.Build.On {
|
||||
if on.Module == "" || on.Module == name || !inSet[on.Module] {
|
||||
continue
|
||||
}
|
||||
for _, base := range entries {
|
||||
if base.Manifest.Module != name && inSet[base.Manifest.Module] && standsOnModule(on, base.Manifest.Module) {
|
||||
if base.Manifest.Module == on.Module {
|
||||
place(base, seen)
|
||||
}
|
||||
}
|
||||
@@ -334,20 +336,10 @@ func standsOn(entries []inventory.Entry, module string) bool {
|
||||
continue
|
||||
}
|
||||
for _, on := range e.Manifest.Build.On {
|
||||
if standsOnModule(on, module) {
|
||||
if on.Module == module {
|
||||
return true
|
||||
}
|
||||
}
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
// standsOnModule is whether a base names the module: as written in a manifest (`module`), or as
|
||||
// recorded after a build, when the mesh has replaced it with the artifact it resolved to
|
||||
// (`artifact-store://<module>/<artifact>@…`). A recorded manifest is what the catalogue holds.
|
||||
func standsOnModule(on catalogue.BuildsOn, module string) bool {
|
||||
if on.Module == module {
|
||||
return true
|
||||
}
|
||||
return strings.HasPrefix(on.Image, "artifact-store://"+module+"/")
|
||||
}
|
||||
|
||||
@@ -4,7 +4,7 @@ go 1.26.0
|
||||
|
||||
require (
|
||||
github.com/jackc/pgx/v5 v5.10.0
|
||||
github.com/nats-io/nats.go v1.54.0
|
||||
github.com/rabbitmq/amqp091-go v1.14.0
|
||||
golang.org/x/crypto v0.57.0
|
||||
)
|
||||
|
||||
@@ -13,6 +13,7 @@ require (
|
||||
github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761 // indirect
|
||||
github.com/jackc/puddle/v2 v2.2.2 // indirect
|
||||
github.com/klauspost/compress v1.20.0 // indirect
|
||||
github.com/nats-io/nats.go v1.54.0 // indirect
|
||||
github.com/nats-io/nkeys v0.4.16 // indirect
|
||||
github.com/nats-io/nuid v1.0.1 // indirect
|
||||
golang.org/x/net v0.58.0 // indirect
|
||||
|
||||
@@ -19,11 +19,15 @@ github.com/nats-io/nuid v1.0.1 h1:5iA8DT8V7q8WK2EScv2padNa/rTESc1KdnPw4TC2paw=
|
||||
github.com/nats-io/nuid v1.0.1/go.mod h1:19wcPz3Ph3q0Jbyiqsd0kePYG7A95tJPxeL+1OSON2c=
|
||||
github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM=
|
||||
github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4=
|
||||
github.com/rabbitmq/amqp091-go v1.14.0 h1:RSaT7aOKt/OrkVUyswPDW29lnRz9psuGmfZFBmLqLek=
|
||||
github.com/rabbitmq/amqp091-go v1.14.0/go.mod h1:Hy4jKW5kQART1u+JkDTF9YYOQUHXqMuhrgxOEeS7G4o=
|
||||
github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME=
|
||||
github.com/stretchr/testify v1.3.0/go.mod h1:M5WIy9Dh21IEIfnGCwXGc5bZfKNJtfHm1UVUgZn+9EI=
|
||||
github.com/stretchr/testify v1.7.0/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg=
|
||||
github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu7U=
|
||||
github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U=
|
||||
go.uber.org/goleak v1.3.0 h1:2K3zAYmnTNqV73imy9J1T3WC+gmCePx2hEGkimedGto=
|
||||
go.uber.org/goleak v1.3.0/go.mod h1:CoHD4mav9JJNrW/WLlf7HGZPjdw8EucARQHekz1X6bE=
|
||||
golang.org/x/crypto v0.57.0 h1:3ZVCjf8Ggz7zneR/EHRVx68Ctf+2pmIMP2UFhh9cC6M=
|
||||
golang.org/x/crypto v0.57.0/go.mod h1:Fdz0i5U6CoizGwLda9DttjSk6qlZo25zYNtR+ycvuZA=
|
||||
golang.org/x/net v0.58.0 h1:ynWG7rqYi4ccpTEuPZ2QGWHktVEM9DMCj9yzDE0Q7To=
|
||||
|
||||
@@ -0,0 +1,67 @@
|
||||
package broker_test
|
||||
|
||||
import (
|
||||
"regexp"
|
||||
"testing"
|
||||
|
||||
"github.com/novox/mesh-controller/internal/broker"
|
||||
"github.com/novox/mesh-controller/internal/link"
|
||||
)
|
||||
|
||||
// The names are written twice, so a test keeps them agreeing.
|
||||
//
|
||||
// `link` imports `broker`, so `broker` cannot import `link` — the queue and exchange names
|
||||
// therefore exist in both. A scoped account naming a queue nothing publishes to produces a
|
||||
// builder that takes no work and says nothing about why, which is the worst kind of silence.
|
||||
//
|
||||
// An external test package, because it may import both without either importing the other.
|
||||
func TestTheNamesTheBrokerScopesAreTheNamesTheLinkUses(t *testing.T) {
|
||||
for _, agreed := range []struct {
|
||||
what string
|
||||
scoped string
|
||||
actually string
|
||||
}{
|
||||
{"the build queue", broker.BuildQueueName, link.BuildQueue},
|
||||
{"the exchange", broker.ExchangeName, link.Exchange},
|
||||
{"a node's queue", broker.QueueFor("somewhere"), link.QueueFor("somewhere")},
|
||||
} {
|
||||
if agreed.scoped != agreed.actually {
|
||||
t.Errorf("%s: the broker scopes %q and the link uses %q",
|
||||
agreed.what, agreed.scoped, agreed.actually)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestABuilderMayWriteToAReplyQueueAndReadNoNodesDeclarations(t *testing.T) {
|
||||
// The scoping, checked as patterns rather than by connecting: what it may write must include
|
||||
// the queue an asker actually waits on, and what it may read must not include any node's.
|
||||
//
|
||||
// This exists because the first version scoped writes to `amq.gen-*` — the name one broker
|
||||
// happens to generate — and the builder built, could not answer, and the connection simply
|
||||
// closed saying only "not allowed to publish to exchange \'\'". Answering through the
|
||||
// exchange is what removed the need for any of that.
|
||||
write := regexp.MustCompile("^" + regexp.QuoteMeta(broker.ExchangeName) + "$")
|
||||
if !write.MatchString(broker.ExchangeName) {
|
||||
t.Error("a builder may not write to the exchange, so it can take work and never answer")
|
||||
}
|
||||
// And not the default exchange, where permission is per exchange rather than per queue — a
|
||||
// builder allowed to use it could publish into any node's queue.
|
||||
//
|
||||
// Confirmed against a real broker as well, and worth recording how that nearly went wrong:
|
||||
// an unconfirmed publish is asynchronous, so a refusal arrives as a channel close afterwards
|
||||
// and a naive check reports success. With publisher confirms the broker's refusal is
|
||||
// immediate. **A negative security assertion made against an asynchronous call is not an
|
||||
// assertion.**
|
||||
if write.MatchString("") {
|
||||
t.Error("a builder may publish to the default exchange, and so into any node's queue")
|
||||
}
|
||||
|
||||
read := regexp.MustCompile("^" + regexp.QuoteMeta(broker.BuildQueueName) + "$")
|
||||
if !read.MatchString(broker.BuildQueueName) {
|
||||
t.Error("a builder may not read the build queue")
|
||||
}
|
||||
if read.MatchString(link.QueueFor("someone-else")) {
|
||||
// A build machine is not a node, and a node's queue carries its declarations.
|
||||
t.Error("a builder may read another machine's declarations")
|
||||
}
|
||||
}
|
||||
@@ -194,3 +194,16 @@ func TestTheAddressPortFollowsThePortTwin(t *testing.T) {
|
||||
t.Fatalf("the address is %q; the node put the bus on 5679", b.Address)
|
||||
}
|
||||
}
|
||||
|
||||
func TestTheManagementPortFollowsThePortTwin(t *testing.T) {
|
||||
t.Setenv(ManagementVar, "http://guest:guest@127.0.0.1:15672")
|
||||
t.Setenv(ManagementVar+"_FILE", "")
|
||||
t.Setenv(ManagementVar+"_PORT", "15673")
|
||||
m, err := ManagementFromEnvironment()
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if m.base.Host != "127.0.0.1:15673" {
|
||||
t.Fatalf("the management API is at %q; the node put it on 15673", m.base.Host)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,366 @@
|
||||
package broker
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"github.com/novox/mesh-controller/internal/envfile"
|
||||
"io"
|
||||
"net/http"
|
||||
"net/url"
|
||||
"regexp"
|
||||
"strings"
|
||||
"time"
|
||||
)
|
||||
|
||||
// The mesh runs the broker, so there is no chicken-and-egg in a node needing an account before it
|
||||
// can connect: the account is created when the token is issued, and the one-time secret in that
|
||||
// token IS the password. A node's first connection is already authenticated, and enrolment is
|
||||
// what happens over it.
|
||||
//
|
||||
// novox/hq ADR 0004's *a node holds its own identity and nothing else* is why the account is per
|
||||
// node rather than shared. A shared enrolment account would let any node consume another's queue,
|
||||
// which is the shared-credential fault that record exists to remove, reappearing at the transport.
|
||||
|
||||
// ManagementVar holds the broker's management API, credentials included.
|
||||
const ManagementVar = "MESH_BROKER_MANAGEMENT"
|
||||
|
||||
// safeName is what a node may be called at the broker.
|
||||
//
|
||||
// The name goes into a URL path and into permission patterns, which are regular expressions. A
|
||||
// name carrying a `.` or a `*` would silently widen what that node may reach — so it is
|
||||
// constrained here rather than escaped later, because an escape that is forgotten once is a node
|
||||
// reading everybody's queues.
|
||||
var safeName = regexp.MustCompile(`^[a-z0-9][a-z0-9-]{0,62}$`)
|
||||
|
||||
// Management is the broker's administrative interface.
|
||||
type Management struct {
|
||||
base *url.URL
|
||||
client *http.Client
|
||||
}
|
||||
|
||||
// ManagementFromEnvironment reads where the management API is, if it is configured.
|
||||
//
|
||||
// On the port MESH_BROKER_MANAGEMENT_PORT names when the node moved it (novox/hq 04-ISSUES/102);
|
||||
// the URL's own port otherwise.
|
||||
func ManagementFromEnvironment() (*Management, error) {
|
||||
raw, err := envfile.Placed(ManagementVar)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if raw == "" {
|
||||
return nil, ErrNotConfigured
|
||||
}
|
||||
base, err := url.Parse(raw)
|
||||
if err != nil || base.Host == "" {
|
||||
// The value carries a password, so it is not quoted back.
|
||||
return nil, fmt.Errorf("%s is not a usable URL", ManagementVar)
|
||||
}
|
||||
return &Management{base: base, client: &http.Client{Timeout: 15 * time.Second}}, nil
|
||||
}
|
||||
|
||||
// QueueFor is the queue a node consumes from. One per node, named after it.
|
||||
func QueueFor(node string) string { return "node." + node }
|
||||
|
||||
// ExchangeName is where nodes publish what they have to say. One exchange, and the control plane
|
||||
// is the only consumer behind it (novox/hq ADR 0006 — one consumer, so two cannot silently split
|
||||
// the traffic between them).
|
||||
const ExchangeName = "mesh"
|
||||
|
||||
// BuildQueueName is where build work waits. Duplicated from `link` rather than imported, for the
|
||||
// same reason QueueFor above is: this package must not depend on the one that uses it, and a
|
||||
// constant that differed would be caught by the test that asserts they agree.
|
||||
const BuildQueueName = "builds"
|
||||
|
||||
// The events bus (novox/hq ADR 0042): one topic exchange every event rides, a second for tool
|
||||
// RPC kept apart, and a dead-letter home for a poison event. The foundation owns these — a module's
|
||||
// account cannot declare them, only bind its own queue to the events one.
|
||||
const (
|
||||
EventsExchangeName = "mesh.events"
|
||||
RPCExchangeName = "mesh.rpc"
|
||||
DeadExchangeName = "mesh.events.dead"
|
||||
)
|
||||
|
||||
// ModuleQueueFor is the durable queue a module consumes its events from — one per module per node
|
||||
// (novox/hq ADR 0042), named so the account that may read it is exactly this module's.
|
||||
func ModuleQueueFor(node, module string) string { return node + "." + module + ".events" }
|
||||
|
||||
// modulePermissions is a module's authority on the bus, derived from its manifest (novox/hq
|
||||
// ADR 0043): what it consumes and what it emits, and nothing else. Pure, so the scope is tested as
|
||||
// patterns without a broker — the way a builder's is.
|
||||
//
|
||||
// A note on the limit: the broker's write permission is per exchange, not per routing key (LavinMQ
|
||||
// has no topic permissions), so an emitting module is granted the events exchange whole. ADR 0042's
|
||||
// origin reservation — a module publishes only under `module.<self>.*` — is stamped by the sdk, not
|
||||
// enforced here; that gap is the broker's, and is recorded rather than hidden. A pure consumer like
|
||||
// the audit logger is unaffected: it is granted no write to the exchange at all.
|
||||
func modulePermissions(node, module string, emits, consumes []string) (configure, write, read string) {
|
||||
queue := regexp.QuoteMeta(ModuleQueueFor(node, module))
|
||||
events := regexp.QuoteMeta(EventsExchangeName)
|
||||
rpc := regexp.QuoteMeta(RPCExchangeName)
|
||||
// A module serves each of its tools on its own queue, namespaced by the module (novox/hq
|
||||
// ADR 0047) — serve.<module>.<tool> — so the account may declare, bind and read exactly its own,
|
||||
// and no other module's.
|
||||
serve := "serve\\." + regexp.QuoteMeta(module) + "\\..*"
|
||||
|
||||
// Declare its own events queue and its own tool serve queues.
|
||||
configure = "^(" + queue + "|" + serve + ")$"
|
||||
|
||||
// Write to bind its queue and serve queues (binding is a write on the queue), and to the RPC
|
||||
// exchange to publish replies (ADR 0047: replies ride mesh.rpc, never the default exchange, which
|
||||
// would let it publish into any queue). To the events exchange only if it emits.
|
||||
writes := []string{queue, serve, rpc}
|
||||
if len(emits) > 0 {
|
||||
writes = append(writes, events)
|
||||
}
|
||||
write = "^(" + strings.Join(writes, "|") + ")$"
|
||||
|
||||
// Read its own queue and serve queues to consume them, and the RPC exchange to bind its serve
|
||||
// queues onto. The events exchange to bind onto only if it consumes.
|
||||
reads := []string{queue, serve, rpc}
|
||||
if len(consumes) > 0 {
|
||||
reads = append(reads, events)
|
||||
}
|
||||
read = "^(" + strings.Join(reads, "|") + ")$"
|
||||
return configure, write, read
|
||||
}
|
||||
|
||||
// CreateModuleAccount gives an assigned module its own broker account, scoped by what it emits and
|
||||
// consumes (novox/hq ADR 0043). The account name carries the node so the same module on two machines
|
||||
// holds two accounts, each sealed to its own; the permissions carry the module so one module cannot
|
||||
// read another's queue. Generic — the builder is one instance of this rule, not a separate kind.
|
||||
func (m *Management) CreateModuleAccount(ctx context.Context, node, module, password string, emits, consumes []string) (string, error) {
|
||||
if !safeName.MatchString(node) {
|
||||
return "", fmt.Errorf("%q cannot be part of a broker account: it is a permission pattern", node)
|
||||
}
|
||||
if !safeName.MatchString(module) {
|
||||
return "", fmt.Errorf("%q cannot be part of a broker account: it is a permission pattern", module)
|
||||
}
|
||||
account := node + "-" + module
|
||||
if !safeName.MatchString(account) {
|
||||
return "", fmt.Errorf("%q is not a usable broker account name", account)
|
||||
}
|
||||
|
||||
if err := m.put(ctx, "/api/users/"+url.PathEscape(account),
|
||||
map[string]string{"password": password, "tags": ""}); err != nil {
|
||||
return "", fmt.Errorf("cannot create the broker account for %s on %s: %w", module, node, err)
|
||||
}
|
||||
|
||||
configure, write, read := modulePermissions(node, module, emits, consumes)
|
||||
if err := m.put(ctx, "/api/permissions/%2f/"+url.PathEscape(account), map[string]string{
|
||||
"configure": configure, "write": write, "read": read,
|
||||
}); err != nil {
|
||||
return "", fmt.Errorf("cannot scope the broker account for %s on %s: %w", module, node, err)
|
||||
}
|
||||
return account, nil
|
||||
}
|
||||
|
||||
// EnsureModuleQueue declares a consuming module's queue with its dead-letter exchange, idempotently.
|
||||
// The foundation declares it because a scoped module account may not: the broker refuses a queue with
|
||||
// a dead-letter exchange to a non-administrator (novox/hq ADR 0043), so a consumer passively checks
|
||||
// the queue the mesh made rather than declaring its own.
|
||||
func (m *Management) EnsureModuleQueue(ctx context.Context, node, module string) error {
|
||||
queue := ModuleQueueFor(node, module)
|
||||
if err := m.put(ctx, "/api/queues/%2f/"+url.PathEscape(queue), map[string]any{
|
||||
"durable": true,
|
||||
"arguments": map[string]any{"x-dead-letter-exchange": DeadExchangeName},
|
||||
}); err != nil {
|
||||
return fmt.Errorf("cannot declare the queue for %s on %s: %w", module, node, err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// EnsureEventExchanges declares the bus's exchanges and the dead-letter home, idempotently. The
|
||||
// foundation owns them (a module's account may not declare an exchange), and a dead-letter exchange
|
||||
// with no queue behind it drops what it receives — so a durable queue bound to `#` retains a poison
|
||||
// event for inspection, which is the whole reason the trail exists.
|
||||
func (m *Management) EnsureEventExchanges(ctx context.Context) error {
|
||||
for _, exchange := range []string{EventsExchangeName, RPCExchangeName, DeadExchangeName} {
|
||||
if err := m.put(ctx, "/api/exchanges/%2f/"+url.PathEscape(exchange),
|
||||
map[string]any{"type": "topic", "durable": true}); err != nil {
|
||||
return fmt.Errorf("cannot declare the %s exchange: %w", exchange, err)
|
||||
}
|
||||
}
|
||||
if err := m.put(ctx, "/api/queues/%2f/"+url.PathEscape(DeadExchangeName),
|
||||
map[string]any{"durable": true}); err != nil {
|
||||
return fmt.Errorf("cannot declare the dead-letter queue: %w", err)
|
||||
}
|
||||
if err := m.post(ctx, "/api/bindings/%2f/e/"+url.PathEscape(DeadExchangeName)+
|
||||
"/q/"+url.PathEscape(DeadExchangeName), map[string]string{"routing_key": "#"}); err != nil {
|
||||
return fmt.Errorf("cannot bind the dead-letter queue: %w", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// CreateNodeAccount gives a node its own broker account, with the token's secret as the password.
|
||||
//
|
||||
// Scoped so a node can reach its own queue and the one exchange, and nothing else. The patterns
|
||||
// are anchored: a node called `laptop` must not be able to read `laptop-of-somebody-else`.
|
||||
func (m *Management) CreateNodeAccount(ctx context.Context, node, password string) error {
|
||||
if !safeName.MatchString(node) {
|
||||
return fmt.Errorf(
|
||||
"%q cannot be a broker account name: it becomes part of a permission pattern, so it "+
|
||||
"is lower-case letters, digits and dashes", node)
|
||||
}
|
||||
|
||||
if err := m.put(ctx, "/api/users/"+url.PathEscape(node),
|
||||
map[string]string{"password": password, "tags": ""}); err != nil {
|
||||
return fmt.Errorf("cannot create the broker account for %s: %w", node, err)
|
||||
}
|
||||
|
||||
queue := regexp.QuoteMeta(QueueFor(node))
|
||||
if err := m.put(ctx, "/api/permissions/%2f/"+url.PathEscape(node), map[string]string{
|
||||
"configure": "^" + queue + "$",
|
||||
"write": "^(" + regexp.QuoteMeta(ExchangeName) + "|" + queue + ")$",
|
||||
"read": "^" + queue + "$",
|
||||
}); err != nil {
|
||||
return fmt.Errorf("cannot scope the broker account for %s: %w", node, err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// CreateBuilderAccount scopes an account to taking build work and answering it.
|
||||
//
|
||||
// **A build machine is not a node**, and giving it a node's account would let it read another
|
||||
// machine's declarations. What it needs is narrower and different: read the build queue, and
|
||||
// write to the exchange and to whatever temporary queue an asker is waiting on.
|
||||
//
|
||||
// The reply queues are the reason `write` is not simply the exchange. `RequestBuild` declares an
|
||||
// exclusive queue with a generated name and waits on it, so a builder that could not write to it
|
||||
// could take work and never answer — which is the failure that looks like a builder that is not
|
||||
// running.
|
||||
func (m *Management) CreateBuilderAccount(ctx context.Context, name, password string) error {
|
||||
if !safeName.MatchString(name) {
|
||||
return fmt.Errorf(
|
||||
"%q cannot be a broker account name: it becomes part of a permission pattern, so it "+
|
||||
"is lower-case letters, digits and dashes", name)
|
||||
}
|
||||
|
||||
if err := m.put(ctx, "/api/users/"+url.PathEscape(name),
|
||||
map[string]string{"password": password, "tags": ""}); err != nil {
|
||||
return fmt.Errorf("cannot create the broker account for %s: %w", name, err)
|
||||
}
|
||||
|
||||
builds := regexp.QuoteMeta(BuildQueueName)
|
||||
if err := m.put(ctx, "/api/permissions/%2f/"+url.PathEscape(name), map[string]string{
|
||||
// It declares the build queue, because whichever builder starts first must be able to —
|
||||
// and a queue nobody may declare is a queue that exists only if the control plane has
|
||||
// already run, which makes the order they start in matter.
|
||||
"configure": "^" + builds + "$",
|
||||
// Two exchanges, and nothing else. **Not the default exchange**: permission there is
|
||||
// granted per exchange rather than per queue, so a builder allowed to use it could
|
||||
// publish into any node's queue — the privilege a build machine most obviously should
|
||||
// not have. Answers go through the node exchange; announcing what was built goes through
|
||||
// the events exchange, which is a different act with a different audience (novox/hq
|
||||
// ADR 0072). A builder that could answer and not announce would leave the module graph
|
||||
// knowing less than the registry does.
|
||||
"write": "^(" + regexp.QuoteMeta(ExchangeName) + "|" + regexp.QuoteMeta(EventsExchangeName) + ")$",
|
||||
// The build queue and nothing else. Not another machine's declarations.
|
||||
"read": "^" + builds + "$",
|
||||
}); err != nil {
|
||||
return fmt.Errorf("cannot scope the broker account for %s: %w", name, err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// RemoveNodeAccount withdraws a node's access.
|
||||
func (m *Management) RemoveNodeAccount(ctx context.Context, node string) error {
|
||||
if !safeName.MatchString(node) {
|
||||
return fmt.Errorf("%q is not a broker account name", node)
|
||||
}
|
||||
return m.do(ctx, http.MethodDelete, "/api/users/"+url.PathEscape(node), nil)
|
||||
}
|
||||
|
||||
// Accounts lists the broker's users, so a picture can be read from the system rather than assumed
|
||||
// (novox/hq ADR 0018).
|
||||
func (m *Management) Accounts(ctx context.Context) ([]string, error) {
|
||||
body, err := m.get(ctx, "/api/users")
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
var users []struct {
|
||||
Name string `json:"name"`
|
||||
}
|
||||
if err := json.Unmarshal(body, &users); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
names := make([]string, 0, len(users))
|
||||
for _, u := range users {
|
||||
names = append(names, u.Name)
|
||||
}
|
||||
return names, nil
|
||||
}
|
||||
|
||||
func (m *Management) put(ctx context.Context, path string, body any) error {
|
||||
return m.do(ctx, http.MethodPut, path, body)
|
||||
}
|
||||
|
||||
func (m *Management) post(ctx context.Context, path string, body any) error {
|
||||
return m.do(ctx, http.MethodPost, path, body)
|
||||
}
|
||||
|
||||
func (m *Management) get(ctx context.Context, path string) ([]byte, error) {
|
||||
request, err := m.request(ctx, http.MethodGet, path, nil)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
response, err := m.client.Do(request)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
defer response.Body.Close()
|
||||
if response.StatusCode >= 300 {
|
||||
return nil, fmt.Errorf("the broker's management API answered %s to GET %s",
|
||||
response.Status, path)
|
||||
}
|
||||
return io.ReadAll(io.LimitReader(response.Body, 1<<20))
|
||||
}
|
||||
|
||||
func (m *Management) do(ctx context.Context, method, path string, body any) error {
|
||||
request, err := m.request(ctx, method, path, body)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
response, err := m.client.Do(request)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer response.Body.Close()
|
||||
if response.StatusCode >= 300 {
|
||||
detail, _ := io.ReadAll(io.LimitReader(response.Body, 4096))
|
||||
return fmt.Errorf("the broker's management API answered %s to %s %s: %s",
|
||||
response.Status, method, path, strings.TrimSpace(string(detail)))
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (m *Management) request(ctx context.Context, method, path string, body any) (*http.Request, error) {
|
||||
var payload io.Reader
|
||||
if body != nil {
|
||||
raw, err := json.Marshal(body)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
payload = bytes.NewReader(raw)
|
||||
}
|
||||
|
||||
// Path joined by hand rather than through url.Parse: %2f is the default vhost and must reach
|
||||
// the broker still encoded. Parsing would decode it to a slash and address a different route.
|
||||
target := strings.TrimSuffix(m.base.String(), "/")
|
||||
if user := m.base.User; user != nil {
|
||||
target = strings.TrimSuffix(m.base.Scheme+"://"+m.base.Host, "/")
|
||||
}
|
||||
request, err := http.NewRequestWithContext(ctx, method, target+path, payload)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if user := m.base.User; user != nil {
|
||||
password, _ := user.Password()
|
||||
request.SetBasicAuth(user.Username(), password)
|
||||
}
|
||||
if body != nil {
|
||||
request.Header.Set("Content-Type", "application/json")
|
||||
}
|
||||
return request, nil
|
||||
}
|
||||
@@ -0,0 +1,96 @@
|
||||
package broker_test
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"regexp"
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
"github.com/novox/mesh-controller/internal/broker"
|
||||
)
|
||||
|
||||
// The audit logger consumes everything and emits nothing. Its account must let it declare and read
|
||||
// its own queue and read the events exchange to bind onto — and must not reach another module's
|
||||
// queue, nor grant any write to the events exchange (novox/hq ADR 0043).
|
||||
func TestAConsumerReadsItsOwnQueueAndTheEventsExchangeAndNoOthers(t *testing.T) {
|
||||
// Rebuilt through CreateModuleAccount's own path by asking for the same scope it would apply.
|
||||
// The queue this module reads:
|
||||
mine := broker.ModuleQueueFor("anchor", "audit-logger")
|
||||
other := broker.ModuleQueueFor("anchor", "plex")
|
||||
|
||||
read := scope(t, "read", "anchor", "audit-logger", nil, []string{"#"})
|
||||
if !read.MatchString(mine) {
|
||||
t.Error("the audit logger may not read its own queue, so it consumes nothing")
|
||||
}
|
||||
if !read.MatchString(broker.EventsExchangeName) {
|
||||
t.Error("the audit logger may not read the events exchange, so it cannot bind onto it")
|
||||
}
|
||||
if read.MatchString(other) {
|
||||
t.Error("the audit logger may read another module's queue")
|
||||
}
|
||||
|
||||
// It emits nothing, so it is granted no write to the events exchange — only its own queue, to bind.
|
||||
write := scope(t, "write", "anchor", "audit-logger", nil, []string{"#"})
|
||||
if write.MatchString(broker.EventsExchangeName) {
|
||||
t.Error("a pure consumer was granted write to the events exchange")
|
||||
}
|
||||
if !write.MatchString(mine) {
|
||||
t.Error("the audit logger may not write to its own queue, so it cannot bind it")
|
||||
}
|
||||
}
|
||||
|
||||
// An emitter is granted the events exchange to write; a consumer is not.
|
||||
func TestAnEmitterMayWriteTheEventsExchangeAndAConsumerMayNot(t *testing.T) {
|
||||
emitter := scope(t, "write", "anchor", "umami", []string{"module.umami.site.created"}, nil)
|
||||
if !emitter.MatchString(broker.EventsExchangeName) {
|
||||
t.Error("an emitting module may not write the events exchange, so it cannot emit")
|
||||
}
|
||||
}
|
||||
|
||||
// scope reconstructs one of the three permission patterns CreateModuleAccount would apply, by
|
||||
// reading it back from a captured request against a stub management API.
|
||||
func scope(t *testing.T, which, node, module string, emits, consumes []string) *regexp.Regexp {
|
||||
t.Helper()
|
||||
pat := capturePermission(t, which, node, module, emits, consumes)
|
||||
re, err := regexp.Compile(pat)
|
||||
if err != nil {
|
||||
t.Fatalf("the %s pattern does not compile: %v", which, err)
|
||||
}
|
||||
return re
|
||||
}
|
||||
|
||||
// capturePermission runs CreateModuleAccount against a stub management API and returns the pattern
|
||||
// it set for `which` ("configure"/"write"/"read"). The scope is tested where it is applied, not
|
||||
// reconstructed by the test — so a change to the mapping cannot pass a test that hard-codes the old
|
||||
// one.
|
||||
func capturePermission(t *testing.T, which, node, module string, emits, consumes []string) string {
|
||||
t.Helper()
|
||||
var captured map[string]string
|
||||
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
if strings.HasPrefix(r.URL.Path, "/api/permissions/") {
|
||||
_ = json.NewDecoder(r.Body).Decode(&captured)
|
||||
}
|
||||
w.WriteHeader(http.StatusNoContent)
|
||||
}))
|
||||
defer server.Close()
|
||||
|
||||
t.Setenv(broker.ManagementVar, server.URL)
|
||||
m, err := broker.ManagementFromEnvironment()
|
||||
if err != nil {
|
||||
t.Fatalf("stub management not usable: %v", err)
|
||||
}
|
||||
if _, err := m.CreateModuleAccount(context.Background(), node, module, "pw", emits, consumes); err != nil {
|
||||
t.Fatalf("CreateModuleAccount: %v", err)
|
||||
}
|
||||
if captured == nil {
|
||||
t.Fatal("no permissions were set")
|
||||
}
|
||||
pattern, ok := captured[which]
|
||||
if !ok {
|
||||
t.Fatalf("no %s permission was set; got %v", which, captured)
|
||||
}
|
||||
return pattern
|
||||
}
|
||||
+19
-11
@@ -68,16 +68,24 @@ func BareAddress(address string) string {
|
||||
return "nats://" + bare
|
||||
}
|
||||
|
||||
// BusAddress is where the mesh's bus is, with this process's credential, and refuses to be empty:
|
||||
// there is one bus, and a control plane without it can hold records and answer nothing (novox/hq
|
||||
// ADR 0131, design 28 task 5.5).
|
||||
func BusAddress() (string, error) {
|
||||
address, _, err := OnNATS()
|
||||
if err != nil {
|
||||
return "", err
|
||||
// MustBeOneBus refuses a configuration that names both buses for the mesh's own traffic.
|
||||
//
|
||||
// **Both clients ship and that is the point; both being live is not.** The rollout moves every node
|
||||
// at once (ADR 0116 step 5): a mesh half on each is one where a declaration goes out on one bus and
|
||||
// the report comes back on the other, and nothing anywhere says so — every component would log
|
||||
// success. Refused at start, where it can be said in one sentence.
|
||||
func MustBeOneBus(amqp, nats string) error {
|
||||
if strings.TrimSpace(amqp) != "" && strings.TrimSpace(nats) != "" {
|
||||
return fmt.Errorf(
|
||||
"this control plane is told about both buses (%s and %s) and can only be on one. A mesh "+
|
||||
"half on each is one where a declaration goes out on one and the report comes back "+
|
||||
"on the other, and every component reports success while it happens. The rollout "+
|
||||
"moves every node at once: unset %s to stay, or unset %s to move",
|
||||
AMQPVarName, NATSVar, NATSVar, AMQPVarName)
|
||||
}
|
||||
if address == "" {
|
||||
return "", fmt.Errorf("this control plane has no %s, so it cannot reach the mesh's bus", NATSVar)
|
||||
}
|
||||
return address, nil
|
||||
return nil
|
||||
}
|
||||
|
||||
// AMQPVarName is the variable naming the bus the mesh runs on today. Named here rather than
|
||||
// imported from the link package, for the one direction of dependency.
|
||||
const AMQPVarName = "MESH_BROKER_AMQP"
|
||||
|
||||
@@ -1,11 +1,43 @@
|
||||
package broker
|
||||
|
||||
import (
|
||||
"strings"
|
||||
"testing"
|
||||
)
|
||||
|
||||
// Which bus the mesh is on is one fact, and being told about both is refused.
|
||||
//
|
||||
// **Not a warning.** A mesh half on each bus is one where a declaration goes out on one and the
|
||||
// report comes back on the other, and every component reports success while it happens — which is
|
||||
// the exact failure ADR 0074 exists to catch, arriving through configuration instead of through code.
|
||||
func TestBeingToldAboutBothBusesIsRefused(t *testing.T) {
|
||||
err := MustBeOneBus("amqps://broker:5671/", "nats://bus:4222")
|
||||
if err == nil {
|
||||
t.Fatal("a control plane told about both buses was allowed to start")
|
||||
}
|
||||
// The remedy is in the words, because whoever reads this has to choose one and the wrong choice
|
||||
// is a rollout half done.
|
||||
for _, want := range []string{AMQPVarName, NATSVar, "unset"} {
|
||||
if !strings.Contains(err.Error(), want) {
|
||||
t.Errorf("the refusal does not mention %s: %v", want, err)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// One bus, or none, is ordinary. None is a control plane that publishes nothing and holds records,
|
||||
// which several of its own commands are.
|
||||
func TestOneBusOrNeitherIsAllowed(t *testing.T) {
|
||||
for _, c := range []struct{ what, amqp, nats string }{
|
||||
{"the bus the mesh runs on today", "amqps://broker:5671/", ""},
|
||||
{"the bus being built", "", "nats://bus:4222"},
|
||||
{"neither", "", ""},
|
||||
{"neither, with whitespace for an address", " ", "\t"},
|
||||
} {
|
||||
if err := MustBeOneBus(c.amqp, c.nats); err != nil {
|
||||
t.Errorf("%s was refused: %v", c.what, err)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// The controller's own credential arrives in its address, and has to be readable out of it — its user
|
||||
// is created by the installer at a bootstrap password, before the controller exists to mint one.
|
||||
|
||||
+62
-24
@@ -3,13 +3,16 @@ package link
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"time"
|
||||
|
||||
"github.com/nats-io/nats.go"
|
||||
amqp "github.com/rabbitmq/amqp091-go"
|
||||
)
|
||||
|
||||
// RPCExchange is where a module's tools are asked over the broker, keyed `<module>.<tool>`, and
|
||||
// where the answer comes back, keyed by the asker's reply queue (novox/hq ADR 0047).
|
||||
const RPCExchange = "mesh.rpc"
|
||||
|
||||
// Answer is what a module's tool replies: one of the two, never both.
|
||||
type Answer struct {
|
||||
Result json.RawMessage `json:"result,omitempty"`
|
||||
@@ -24,35 +27,70 @@ type Answer struct {
|
||||
// grants, and should not. The control plane already holds a connection that may, so a person or
|
||||
// an agent asks through it, and every question passes one process where an audit belongs.
|
||||
//
|
||||
// The answer comes back on the asker's own inbox, which only the asker may read; the serving
|
||||
// module answers there and nowhere else. The bus refuses a request nothing serves at once, so a
|
||||
// module that is down or a tool that does not exist is said now rather than after the whole wait.
|
||||
func Ask(ctx context.Context, bus Bus, module, tool string, args json.RawMessage,
|
||||
// The reply queue is the caller's own, server-named and exclusive, bound to the RPC exchange under
|
||||
// its own name: a serving module answers through that exchange and never the default one, whose
|
||||
// permission is per exchange rather than per queue. The correlation is checked rather than
|
||||
// assumed, as every RPC here is.
|
||||
func Ask(ctx context.Context, channel *amqp.Channel, module, tool string, args json.RawMessage,
|
||||
timeout time.Duration) (Answer, error) {
|
||||
|
||||
if len(args) == 0 {
|
||||
args = json.RawMessage(`{}`)
|
||||
}
|
||||
key := module + "." + tool
|
||||
reply, err := bus.AskTool(ctx, module, tool, args, timeout)
|
||||
replies, err := channel.QueueDeclare("", false, true, true, false, nil)
|
||||
if err != nil {
|
||||
if errors.Is(err, nats.ErrNoResponders) {
|
||||
return Answer{}, fmt.Errorf(
|
||||
"nothing serves %s: no runtime has bound %q on the bus. The module is not "+
|
||||
"assigned, its runtime is not up, or it serves no such tool — `status` says "+
|
||||
"whether the machine carrying it has applied", module, key)
|
||||
}
|
||||
if errors.Is(err, context.DeadlineExceeded) || errors.Is(err, nats.ErrTimeout) {
|
||||
return Answer{}, fmt.Errorf(
|
||||
"%s did not answer within %s. Its runtime serves %q when it is up and has bound "+
|
||||
"the bus — `status` says whether the machine carrying it has applied",
|
||||
module, timeout, key)
|
||||
}
|
||||
return Answer{}, err
|
||||
}
|
||||
if err := channel.QueueBind(replies.Name, replies.Name, RPCExchange, false, nil); err != nil {
|
||||
return Answer{}, fmt.Errorf("cannot bind a reply queue to %s: %w", RPCExchange, err)
|
||||
}
|
||||
answers, err := channel.ConsumeWithContext(ctx, replies.Name, "", true, true, false, false, nil)
|
||||
if err != nil {
|
||||
return Answer{}, err
|
||||
}
|
||||
|
||||
// Mandatory, so a request nothing consumes comes straight back: a module that is down, or a
|
||||
// tool that does not exist, is said at once rather than after the whole wait.
|
||||
returned := channel.NotifyReturn(make(chan amqp.Return, 1))
|
||||
id := fmt.Sprintf("ask-%d", time.Now().UnixNano())
|
||||
key := module + "." + tool
|
||||
if err := channel.PublishWithContext(ctx, RPCExchange, key, true, false, amqp.Publishing{
|
||||
ContentType: "application/json",
|
||||
CorrelationId: id,
|
||||
ReplyTo: replies.Name,
|
||||
Body: args,
|
||||
}); err != nil {
|
||||
return Answer{}, fmt.Errorf("cannot ask %s: %w", key, err)
|
||||
}
|
||||
var answer Answer
|
||||
if err := json.Unmarshal(reply, &answer); err != nil {
|
||||
return Answer{}, fmt.Errorf("%s answered with something unreadable: %w", key, err)
|
||||
|
||||
waiting, cancel := context.WithTimeout(ctx, timeout)
|
||||
defer cancel()
|
||||
for {
|
||||
select {
|
||||
case back := <-returned:
|
||||
if back.CorrelationId == id {
|
||||
return Answer{}, fmt.Errorf(
|
||||
"nothing serves %s: no runtime has bound %q on the broker. The module is not "+
|
||||
"assigned, its runtime is not up, or it serves no such tool — `status` "+
|
||||
"says whether the machine carrying it has applied", module, key)
|
||||
}
|
||||
case <-waiting.Done():
|
||||
return Answer{}, fmt.Errorf(
|
||||
"%s did not answer within %s. Its runtime serves %q when it is up and has bound "+
|
||||
"the broker — `status` says whether the machine carrying it has applied",
|
||||
module, timeout, key)
|
||||
case delivery, ok := <-answers:
|
||||
if !ok {
|
||||
return Answer{}, fmt.Errorf("the connection closed while waiting for %s", key)
|
||||
}
|
||||
if delivery.CorrelationId != id {
|
||||
continue
|
||||
}
|
||||
var answer Answer
|
||||
if err := json.Unmarshal(delivery.Body, &answer); err != nil {
|
||||
return Answer{}, fmt.Errorf("%s answered with something unreadable: %w", key, err)
|
||||
}
|
||||
return answer, nil
|
||||
}
|
||||
}
|
||||
return answer, nil
|
||||
}
|
||||
|
||||
@@ -1,7 +1,12 @@
|
||||
package link
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"time"
|
||||
|
||||
amqp "github.com/rabbitmq/amqp091-go"
|
||||
)
|
||||
|
||||
// Asking a machine to build a module, and hearing what came out.
|
||||
@@ -17,6 +22,21 @@ import (
|
||||
// holds no opinion about what they contain, and a host that also built things would be a host
|
||||
// with a container runtime requirement and a git dependency (novox/hq ADR 0005).
|
||||
|
||||
// BuildQueue is where build requests wait. One queue, so several build machines can share the
|
||||
// work and each request is done exactly once — which is what a queue is for and what a
|
||||
// per-machine routing key would not give.
|
||||
const BuildQueue = "builds"
|
||||
|
||||
// KeyBuilt is what a builder publishes when it has finished, successfully or not.
|
||||
const KeyBuilt = "built"
|
||||
|
||||
// ReplyQueue is where the answer to one request goes.
|
||||
//
|
||||
// **Named here rather than left to the broker**, so it can be scoped and reasoned about. A broker
|
||||
// generates its own name for an unnamed queue, and a builder permitted to write to whatever that
|
||||
// convention happens to produce works on one broker and silently cannot answer on another.
|
||||
func ReplyQueue(id string) string { return BuildQueue + ".reply." + id }
|
||||
|
||||
// BuildRequest is one module to build.
|
||||
type BuildRequest struct {
|
||||
// ID correlates the answer with the asking. Not the module name: two builds of one module can
|
||||
@@ -98,3 +118,77 @@ type MadeArtifact struct {
|
||||
Kind string `json:"kind"`
|
||||
Reference string `json:"reference"`
|
||||
}
|
||||
|
||||
// RequestBuild asks for a module to be built and waits for the answer.
|
||||
//
|
||||
// Waiting rather than returning immediately, because the thing a person wants after asking for a
|
||||
// build is to know whether it worked. A build that is dispatched and forgotten needs somewhere to
|
||||
// look afterwards, and there is nowhere yet.
|
||||
func RequestBuild(ctx context.Context, channel *amqp.Channel, request BuildRequest,
|
||||
timeout time.Duration) (BuildResult, error) {
|
||||
|
||||
// Its own queue for the answer, declared before the ask. Consuming from the shared exchange
|
||||
// would mean competing with the control plane's own consumer for a message meant for this
|
||||
// caller — which is the fault this package's own doc comment records having had.
|
||||
replies, err := channel.QueueDeclare(ReplyQueue(request.ID), false, true, true, false, nil)
|
||||
if err != nil {
|
||||
return BuildResult{}, err
|
||||
}
|
||||
// Bound to the exchange, and the answer comes back through it.
|
||||
//
|
||||
// **A builder never publishes to the default exchange**, because permission there is per
|
||||
// exchange and not per queue — a builder allowed to use it could publish into any node's
|
||||
// queue, which is the privilege a build machine most obviously should not have. Found by
|
||||
// running it: the builder built, could not answer, and the connection closed saying only
|
||||
// "not allowed to publish to exchange ''".
|
||||
//
|
||||
// The cost is that every asker sees every result, which is why the correlation is checked
|
||||
// below rather than assumed.
|
||||
if err := channel.QueueBind(replies.Name, KeyBuilt, Exchange, false, nil); err != nil {
|
||||
return BuildResult{}, err
|
||||
}
|
||||
answers, err := channel.ConsumeWithContext(ctx, replies.Name, "", true, true, false, false, nil)
|
||||
if err != nil {
|
||||
return BuildResult{}, err
|
||||
}
|
||||
|
||||
body, err := json.Marshal(request)
|
||||
if err != nil {
|
||||
return BuildResult{}, err
|
||||
}
|
||||
if err := channel.PublishWithContext(ctx, "", BuildQueue, false, false, amqp.Publishing{
|
||||
ContentType: "application/json",
|
||||
DeliveryMode: amqp.Persistent,
|
||||
CorrelationId: request.ID,
|
||||
ReplyTo: replies.Name,
|
||||
Body: body,
|
||||
}); err != nil {
|
||||
return BuildResult{}, err
|
||||
}
|
||||
|
||||
waiting, cancel := context.WithTimeout(ctx, timeout)
|
||||
defer cancel()
|
||||
for {
|
||||
select {
|
||||
case <-waiting.Done():
|
||||
return BuildResult{}, fmt.Errorf(
|
||||
"no builder answered within %s. Something must be consuming %q, and nothing is "+
|
||||
"— or it is building something that takes longer than this",
|
||||
timeout, BuildQueue)
|
||||
case delivery, ok := <-answers:
|
||||
if !ok {
|
||||
return BuildResult{}, fmt.Errorf("the connection closed while waiting for a build")
|
||||
}
|
||||
var result BuildResult
|
||||
if err := json.Unmarshal(delivery.Body, &result); err != nil {
|
||||
return BuildResult{}, fmt.Errorf("a builder answered with something unreadable: %w", err)
|
||||
}
|
||||
if result.ID != request.ID {
|
||||
// Somebody else's answer on this queue. Ignored rather than returned, because
|
||||
// returning it would attribute one build's outcome to another's.
|
||||
continue
|
||||
}
|
||||
return result, nil
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,202 @@
|
||||
package link
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"time"
|
||||
|
||||
amqp "github.com/rabbitmq/amqp091-go"
|
||||
)
|
||||
|
||||
// The build flow on the bus the mesh runs on today.
|
||||
//
|
||||
// Moved behind the seam rather than changed. The queue, the reply binding and the correlation are
|
||||
// what they were, because the mesh is running on this.
|
||||
|
||||
// currentBuilds asks for builds over a channel.
|
||||
type currentBuilds struct{ channel *amqp.Channel }
|
||||
|
||||
// BuildsOverCurrent is the asking side on the bus the mesh has.
|
||||
func BuildsOverCurrent(channel *amqp.Channel) Builders { return currentBuilds{channel: channel} }
|
||||
|
||||
func (b currentBuilds) Close() {}
|
||||
|
||||
func (b currentBuilds) Submit(ctx context.Context, request BuildRequest,
|
||||
wait time.Duration) (BuildResult, error) {
|
||||
|
||||
// Its own queue for the answer, declared before the ask. Consuming from the shared exchange
|
||||
// would mean competing with the controller's own consumer for a message meant for this caller.
|
||||
replies, err := b.channel.QueueDeclare(ReplyQueue(request.ID), false, true, true, false, nil)
|
||||
if err != nil {
|
||||
return BuildResult{}, err
|
||||
}
|
||||
// **A builder never publishes to the default exchange**, because permission there is per
|
||||
// exchange and not per queue — a builder allowed to use it could publish into any node's queue,
|
||||
// which is the privilege a build machine most obviously should not have. The cost is that every
|
||||
// asker sees every result, which is why the correlation is checked below rather than assumed.
|
||||
if err := b.channel.QueueBind(replies.Name, KeyBuilt, Exchange, false, nil); err != nil {
|
||||
return BuildResult{}, err
|
||||
}
|
||||
answers, err := b.channel.ConsumeWithContext(ctx, replies.Name, "", true, true, false, false, nil)
|
||||
if err != nil {
|
||||
return BuildResult{}, err
|
||||
}
|
||||
|
||||
body, err := json.Marshal(request)
|
||||
if err != nil {
|
||||
return BuildResult{}, err
|
||||
}
|
||||
if err := b.channel.PublishWithContext(ctx, "", BuildQueue, false, false, amqp.Publishing{
|
||||
ContentType: "application/json",
|
||||
DeliveryMode: amqp.Persistent,
|
||||
CorrelationId: request.ID,
|
||||
ReplyTo: replies.Name,
|
||||
Body: body,
|
||||
}); err != nil {
|
||||
return BuildResult{}, err
|
||||
}
|
||||
|
||||
waiting, cancel := context.WithTimeout(ctx, wait)
|
||||
defer cancel()
|
||||
for {
|
||||
select {
|
||||
case <-waiting.Done():
|
||||
return BuildResult{}, waitingFor(wait)
|
||||
case delivery, ok := <-answers:
|
||||
if !ok {
|
||||
return BuildResult{}, errors.New("the connection closed while waiting for a build")
|
||||
}
|
||||
result, mine, err := theOutcomeOf(delivery.Body, request.ID)
|
||||
if err != nil {
|
||||
return BuildResult{}, err
|
||||
}
|
||||
if mine {
|
||||
return result, nil
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// --- the machine's side ---------------------------------------------------------------------
|
||||
|
||||
type currentMachine struct {
|
||||
conn *amqp.Connection
|
||||
channel *amqp.Channel
|
||||
on string
|
||||
}
|
||||
|
||||
// MachineOverCurrent takes build work over a channel.
|
||||
func MachineOverCurrent(conn *amqp.Connection, channel *amqp.Channel, on string) BuildMachine {
|
||||
return ¤tMachine{conn: conn, channel: channel, on: on}
|
||||
}
|
||||
|
||||
func (m *currentMachine) Close() {}
|
||||
|
||||
func (m *currentMachine) Take(ctx context.Context, do func(context.Context, Build)) error {
|
||||
if _, err := m.channel.QueueDeclare(BuildQueue, true, false, false, false, nil); err != nil {
|
||||
return err
|
||||
}
|
||||
// One at a time. A machine that took five requests at once would run five container builds
|
||||
// against one runtime and finish all of them slower than it would have finished the first — and
|
||||
// the queue is what shares work between machines, so nothing is lost by it.
|
||||
if err := m.channel.Qos(1, 0, false); err != nil {
|
||||
return err
|
||||
}
|
||||
// Not auto-acknowledged: a request acknowledged on arrival is a build that vanishes if this
|
||||
// process dies mid-way, with nobody waiting on it ever hearing why.
|
||||
requests, err := m.channel.ConsumeWithContext(ctx, BuildQueue, "mesh-builder",
|
||||
false, false, false, false, nil)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
for {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return nil
|
||||
case delivery, ok := <-requests:
|
||||
if !ok {
|
||||
return errors.New("the broker closed the connection")
|
||||
}
|
||||
var request BuildRequest
|
||||
if err := json.Unmarshal(delivery.Body, &request); err != nil {
|
||||
// Unreadable: rejected rather than retried, because the next attempt reads the same
|
||||
// bytes. Nobody waiting hears an answer, which is correct — there was no request.
|
||||
_ = delivery.Reject(false)
|
||||
continue
|
||||
}
|
||||
do(ctx, ¤tBuild{request: request, delivery: delivery, on: m.on, channel: m.channel})
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
type currentBuild struct {
|
||||
request BuildRequest
|
||||
delivery amqp.Delivery
|
||||
on string
|
||||
channel *amqp.Channel
|
||||
}
|
||||
|
||||
func (b *currentBuild) Request() BuildRequest { return b.request }
|
||||
|
||||
// Announce answers and announces, which on this bus are two publishes to two exchanges.
|
||||
//
|
||||
// The reply goes to whoever asked, correlated to their request; the announcement says to the whole
|
||||
// mesh that a module now exists at a commit (novox/hq ADR 0072). Only a successful build is
|
||||
// announced: a failed one produced no module version, and announcing one would put something in the
|
||||
// graph that was never made.
|
||||
func (b *currentBuild) Announce(ctx context.Context, result BuildResult) error {
|
||||
body, err := json.Marshal(result)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if err := b.channel.PublishWithContext(ctx, Exchange, KeyBuilt, false, false, amqp.Publishing{
|
||||
ContentType: "application/json",
|
||||
CorrelationId: result.ID,
|
||||
Body: body,
|
||||
}); err != nil {
|
||||
return fmt.Errorf("cannot answer a build request: %w", err)
|
||||
}
|
||||
if result.Failed != "" || result.Commit == "" {
|
||||
return nil
|
||||
}
|
||||
|
||||
// **Announced under both names on this bus, for exactly as long as this bus lives.**
|
||||
//
|
||||
// A build's outcome belongs to the role now (novox/hq ADR 0121), so a catalogue built from the
|
||||
// current manifests listens for the role's name. A catalogue that is *already running* listens for
|
||||
// the module's, because that is what it was told when it was installed. A rename on a live bus
|
||||
// needs the publisher and the subscriber to change together, and a merge cannot promise that: one
|
||||
// of them is deployed first, and in that window the graph silently stops being updated — which is
|
||||
// the failure this whole change was cleaning up after.
|
||||
//
|
||||
// So both, and the order stops mattering. The module's own name goes with the bus, in step 5's
|
||||
// retirement list; nothing has ever run on the bus being built, so there is no legacy name there
|
||||
// and this doubling has no counterpart.
|
||||
announced := announcementOf(result)
|
||||
if err := EmitEvent(ctx, OverCurrent{Channel: b.channel}, KeyModuleBuilt, "builder", b.on,
|
||||
announced); err != nil {
|
||||
return err
|
||||
}
|
||||
return EmitEvent(ctx, OverCurrent{Channel: b.channel}, KeyRoleBuilt, TheBuildMachine, b.on,
|
||||
announced)
|
||||
}
|
||||
|
||||
func (b *currentBuild) Done() error { return b.delivery.Ack(false) }
|
||||
|
||||
func (b *currentBuild) Hold(time.Duration) error {
|
||||
// No delayed redelivery on this bus: handed back at once, which is what it has always done.
|
||||
return b.delivery.Nack(false, true)
|
||||
}
|
||||
|
||||
// announcementOf is what the mesh is told about a finished build. One function, so the two
|
||||
// transports cannot describe the same build differently.
|
||||
func announcementOf(result BuildResult) map[string]any {
|
||||
return map[string]any{
|
||||
"module": ModuleOf(result.Manifest), "commit": result.Commit,
|
||||
"repository": result.Repository, "path": result.Path, "ref": result.Ref,
|
||||
"manifest": json.RawMessage(result.Manifest), "against": result.Against,
|
||||
"made": result.Made,
|
||||
}
|
||||
}
|
||||
@@ -7,6 +7,7 @@ import (
|
||||
"time"
|
||||
|
||||
"github.com/nats-io/nats.go"
|
||||
amqp "github.com/rabbitmq/amqp091-go"
|
||||
)
|
||||
|
||||
// Bus is what the controller needs of the mesh's bus, **in the mesh's own words rather than a
|
||||
@@ -38,6 +39,51 @@ type Bus interface {
|
||||
|
||||
// --- The bus the mesh runs on today -----------------------------------------------------
|
||||
|
||||
// OverCurrent is the bus the mesh runs on today, until the rollout.
|
||||
type OverCurrent struct{ Channel *amqp.Channel }
|
||||
|
||||
func (b OverCurrent) PublishEvent(ctx context.Context, key, source, node string, body []byte) error {
|
||||
id, err := eventID()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
return b.Channel.PublishWithContext(ctx, EventsExchange, key, false, false, amqp.Publishing{
|
||||
ContentType: "application/json",
|
||||
DeliveryMode: amqp.Persistent,
|
||||
MessageId: id,
|
||||
Timestamp: time.Now().UTC(),
|
||||
Body: body,
|
||||
Headers: amqp.Table{
|
||||
"x-event-id": id,
|
||||
"x-source": source,
|
||||
"x-node": node,
|
||||
"x-time": time.Now().UTC().Format(time.RFC3339),
|
||||
"content-type": "application/json",
|
||||
},
|
||||
})
|
||||
}
|
||||
|
||||
// AskTool is implemented over the existing reply-queue machinery in ask.go; this seam does not
|
||||
// change how it works today.
|
||||
func (b OverCurrent) AskTool(ctx context.Context, module, tool string, args []byte, timeout time.Duration) ([]byte, error) {
|
||||
answer, err := Ask(ctx, b.Channel, module, tool, args, timeout)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return answer.Result, nil
|
||||
}
|
||||
|
||||
func (b OverCurrent) PublishDeclaration(ctx context.Context, node string, body []byte) error {
|
||||
// To the queue directly rather than through an exchange: a declaration is for one node, and
|
||||
// routing it by name through a shared exchange would mean a binding per node that nothing
|
||||
// removes when a node is retired.
|
||||
return b.Channel.PublishWithContext(ctx, "", QueueFor(node), false, false, amqp.Publishing{
|
||||
ContentType: "application/json",
|
||||
DeliveryMode: amqp.Persistent,
|
||||
Body: body,
|
||||
})
|
||||
}
|
||||
|
||||
// --- NATS, the bus being built ----------------------------------------------------------------
|
||||
|
||||
// OverNATS is the bus as a JetStream context.
|
||||
|
||||
@@ -1,16 +1,66 @@
|
||||
package link
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"time"
|
||||
|
||||
"github.com/nats-io/nats.go"
|
||||
|
||||
"github.com/novox/mesh-controller/internal/broker"
|
||||
)
|
||||
|
||||
// OverNats is the controller's outbound on the bus being built: the same three acts the other
|
||||
// transport has, on the subjects the permissions were derived for (design 25). A declaration is a
|
||||
// JetStream publish into NODES, where the node's own consumer waits for it; an event is announced on
|
||||
// the subject its name derives to; a tool is asked by request and reply on the module's tool subject.
|
||||
type OverNats struct{ JS *broker.JetStream }
|
||||
|
||||
// declareSubject is where one node's declaration lands — the NODES stream's subject for it, and the
|
||||
// only subject that node's consumer delivers. The host subscribes exactly this.
|
||||
func declareSubject(node string) string { return "mesh.node." + node + ".declare" }
|
||||
|
||||
func (b OverNats) PublishDeclaration(ctx context.Context, node string, body []byte) error {
|
||||
publish, cancel := context.WithTimeout(ctx, 15*time.Second)
|
||||
defer cancel()
|
||||
if _, err := b.JS.Context().Publish(declareSubject(node), body, nats.Context(publish)); err != nil {
|
||||
return fmt.Errorf("declaring to %s: %w", node, err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (b OverNats) PublishEvent(ctx context.Context, key, source, node string, body []byte) error {
|
||||
// The key is the subject: the controller's own events are named in full, and what a module
|
||||
// emits is derived before it reaches here. Headers carry the envelope the other transport put
|
||||
// in message properties (ADR 0042), so a consumer reads who and when without the payload.
|
||||
msg := nats.NewMsg(key)
|
||||
msg.Data = body
|
||||
msg.Header.Set("x-source", source)
|
||||
msg.Header.Set("x-node", node)
|
||||
msg.Header.Set("x-time", time.Now().UTC().Format(time.RFC3339Nano))
|
||||
if err := b.JS.Conn().PublishMsg(msg); err != nil {
|
||||
return fmt.Errorf("announcing %s: %w", key, err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (b OverNats) AskTool(ctx context.Context, module, tool string, args []byte, timeout time.Duration) ([]byte, error) {
|
||||
ask, cancel := context.WithTimeout(ctx, timeout)
|
||||
defer cancel()
|
||||
reply, err := b.JS.Conn().RequestWithContext(ask, "mesh.mod."+module+".tool."+tool, args)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("asking %s.%s: %w", module, tool, err)
|
||||
}
|
||||
return reply.Data, nil
|
||||
}
|
||||
|
||||
// ConnectNats is Connect for the bus being built: the controller's inbound and outbound over one
|
||||
// JetStream connection the caller has already raised the streams on. Nothing is declared here —
|
||||
// the streams and the controller's consumers are asserted by Raise, before anything is served.
|
||||
func ConnectNats(js *broker.JetStream, enroller Enroller, listener Listener) *Server {
|
||||
return &Server{
|
||||
inbound: Nats(js),
|
||||
bus: OverNATS{JS: js.Context(), Conn: js.Conn()},
|
||||
bus: OverNats{JS: js},
|
||||
js: js,
|
||||
enroller: enroller,
|
||||
listener: listener,
|
||||
|
||||
@@ -24,9 +24,10 @@ import (
|
||||
// — it holds both grants and asks each for its part, which is what the process running them is
|
||||
// for.
|
||||
type Enrolment struct {
|
||||
Inventory *inventory.Inventory
|
||||
Identity *identity.Identity
|
||||
Broker broker.Broker
|
||||
Inventory *inventory.Inventory
|
||||
Identity *identity.Identity
|
||||
Management *broker.Management
|
||||
Broker broker.Broker
|
||||
|
||||
// OnNATS says the mesh's own traffic is on the bus being built, so a node's credential is
|
||||
// minted into the mesh's records and composed into the bus's user list rather than pushed
|
||||
@@ -193,6 +194,17 @@ func (e Enrolment) Enrol(ctx context.Context, request EnrolRequest) (reply Enrol
|
||||
reply.Password = password
|
||||
}
|
||||
|
||||
case e.Management != nil:
|
||||
password, err := freshPassword()
|
||||
if err != nil {
|
||||
return EnrolReply{}, err
|
||||
}
|
||||
if err := e.Management.CreateNodeAccount(ctx, node.Name, password); err != nil {
|
||||
log.Printf("%s is enrolled and its broker password could not be replaced, so it keeps "+
|
||||
"the token's secret as its password: %v", node.Name, err)
|
||||
} else {
|
||||
reply.Password = password
|
||||
}
|
||||
}
|
||||
|
||||
if profile != nil {
|
||||
@@ -249,6 +261,9 @@ func freshPassword() (string, error) {
|
||||
|
||||
var _ Enroller = Enrolment{}
|
||||
|
||||
// ErrNoBrokerManagement is returned when an account cannot be made because nothing was configured.
|
||||
var ErrNoBrokerManagement = errors.New("no broker management configured")
|
||||
|
||||
// Heard records what a node reported about itself.
|
||||
//
|
||||
// A node states; the owning context writes (novox/hq ADR 0006). What a node says it applied is
|
||||
|
||||
@@ -78,6 +78,19 @@ type Control interface {
|
||||
// Took settles the message: acted on, or understood and needing no action.
|
||||
Took() error
|
||||
|
||||
// About names what this message is about — a node's report, one module's move, one build's
|
||||
// outcome — and is said before the store is asked.
|
||||
//
|
||||
// A transport that holds messages **in memory** uses it to set aside anything older it is
|
||||
// holding about the same thing: the older is the past, and letting it come back after the
|
||||
// newer was acted on would undo the newer.
|
||||
//
|
||||
// **This is the one thing holding-in-memory can do that holding-in-the-server cannot**, and
|
||||
// naming it here rather than hiding it is deliberate. On the bus being built the message
|
||||
// belongs to the server and comes back whatever happened meanwhile, so this is ignored and the
|
||||
// digest a report carries answers the same question instead (window.go, design 25 §3).
|
||||
About(what string)
|
||||
|
||||
// Hold keeps the message and asks for it again after the delay — the store window.
|
||||
Hold(after time.Duration) error
|
||||
|
||||
|
||||
@@ -0,0 +1,321 @@
|
||||
package link
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"time"
|
||||
|
||||
amqp "github.com/rabbitmq/amqp091-go"
|
||||
)
|
||||
|
||||
// The consume side on the bus the mesh runs on today.
|
||||
//
|
||||
// Everything here was the serving loop's until the seam went in: the queues, the binds, the
|
||||
// prefetch, and the list of messages the store could not take yet. It moved rather than changed —
|
||||
// the behaviour this transport has is the behaviour it had, because the mesh is running on it and
|
||||
// a bus nothing speaks yet is no reason to alter the one every node is on (ADR 0116).
|
||||
|
||||
// Prefetch is how many messages the bus hands the controller before it has settled them.
|
||||
//
|
||||
// More than one because a message the store could not take is held, unsettled, while the loop
|
||||
// goes on answering others — an enrolment above all, which a host is waiting on (novox/hq issue
|
||||
// 083). Bounded, because what is held is also what the bus has not kept on its own disk as
|
||||
// pending.
|
||||
const Prefetch = 64
|
||||
|
||||
// PrefetchHeadroom is how much of the prefetch is never held, so the loop always has messages to
|
||||
// answer — an enrolment above all — while others wait for the store.
|
||||
const PrefetchHeadroom = 8
|
||||
|
||||
// TryAgainAfter is how often the held are looked at. A store comes back in seconds, and a report a
|
||||
// few seconds late is still current.
|
||||
const TryAgainAfter = 2 * time.Second
|
||||
|
||||
// currentInbound consumes what nodes say over the bus the mesh has.
|
||||
type currentInbound struct {
|
||||
conn *amqp.Connection
|
||||
channel *amqp.Channel
|
||||
// upgrades and catchups are bound only when something is listening (Also).
|
||||
upgrades bool
|
||||
catchups bool
|
||||
// held is every message the store could not take, by the bus's own delivery tag. Kept here
|
||||
// rather than in the serving loop because holding a delivery unacknowledged is this
|
||||
// transport's way of keeping it, and the other's is to hand it back to the server.
|
||||
held map[uint64]*holding
|
||||
// again is how often the held are looked at; zero means TryAgainAfter. Set by tests.
|
||||
again time.Duration
|
||||
}
|
||||
|
||||
// holding is one message kept for the store, and when to try it again.
|
||||
type holding struct {
|
||||
message *currentControl
|
||||
due time.Time
|
||||
about string
|
||||
}
|
||||
|
||||
// Current is the consume side of the bus the mesh runs on today.
|
||||
func Current(conn *amqp.Connection, channel *amqp.Channel) Inbound {
|
||||
return ¤tInbound{conn: conn, channel: channel, held: map[uint64]*holding{}}
|
||||
}
|
||||
|
||||
// Also binds the queue one more kind arrives on.
|
||||
//
|
||||
// The kinds nodes publish all share one queue and are bound at Connect, because a node may
|
||||
// publish any of them and binding one while forgetting another is a message the bus accepts, finds
|
||||
// no queue for, and drops — the publisher sees success and the consumer sees nothing. The two that
|
||||
// are events get their own queue each, and only when something is listening.
|
||||
func (c *currentInbound) Also(kind string) error {
|
||||
switch kind {
|
||||
case KindSourceMoved:
|
||||
// Not followed on the bus the mesh is leaving: the forge's merges are announced on the
|
||||
// new one, and this transport goes with the move (design 28, task 5.5).
|
||||
return nil
|
||||
case KindModuleMoved:
|
||||
if err := c.bindEvent(UpgradeQueue, KeyModuleUpgraded); err != nil {
|
||||
return err
|
||||
}
|
||||
c.upgrades = true
|
||||
case KindCatchUp:
|
||||
if err := c.bindEvent(CatchUpQueue, KeyCatchingUp); err != nil {
|
||||
return err
|
||||
}
|
||||
c.catchups = true
|
||||
default:
|
||||
return fmt.Errorf("nothing binds a queue for %s on this bus", kind)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (c *currentInbound) bindEvent(queue, key string) error {
|
||||
if _, err := c.channel.QueueDeclare(queue, true, false, false, false, nil); err != nil {
|
||||
return fmt.Errorf("cannot declare the %s queue: %w", queue, err)
|
||||
}
|
||||
if err := c.channel.QueueBind(queue, key, EventsExchange, false, nil); err != nil {
|
||||
return fmt.Errorf("cannot bind %s to %s/%s: %w", queue, EventsExchange, key, err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (c *currentInbound) Close() {}
|
||||
|
||||
// Receive consumes until the context ends.
|
||||
//
|
||||
// One consumer per queue, deliberately: with two on one queue the bus would round-robin between
|
||||
// them and each would receive half of what it expects — a fault this project has already had,
|
||||
// between a module's daemon and its capability server.
|
||||
func (c *currentInbound) Receive(ctx context.Context, act func(context.Context, Control)) error {
|
||||
// A bounded prefetch rather than one. The loop still takes messages one at a time; what the
|
||||
// prefetch buys is that a message the store could not take can be held while the loop goes on
|
||||
// to the next, instead of every enrolment waiting behind it (novox/hq issue 083). Anything
|
||||
// held goes back to the bus if the controller stops, because nothing held is acknowledged.
|
||||
if err := c.channel.Qos(Prefetch, 0, false); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
deliveries, err := c.channel.ConsumeWithContext(ctx, ControlQueue, "control-plane",
|
||||
false, false, false, false, nil)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
// Its own queue and its own consumer for each event, for the reason above: two consumers on
|
||||
// one queue split its messages between them, and an upgrade or a catch-up request going to
|
||||
// whichever half was not listening is a gap that looks like a working mesh.
|
||||
var upgrades, catchups <-chan amqp.Delivery
|
||||
if c.upgrades {
|
||||
upgrades, err = c.channel.ConsumeWithContext(ctx, UpgradeQueue, "control-plane-upgrades",
|
||||
false, false, false, false, nil)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
if c.catchups {
|
||||
catchups, err = c.channel.ConsumeWithContext(ctx, CatchUpQueue, "control-plane-catchup",
|
||||
false, false, false, false, nil)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
||||
closed := c.conn.NotifyClose(make(chan *amqp.Error, 1))
|
||||
|
||||
again := c.again
|
||||
if again == 0 {
|
||||
again = TryAgainAfter
|
||||
}
|
||||
ticker := time.NewTicker(again)
|
||||
defer ticker.Stop()
|
||||
|
||||
for {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return nil
|
||||
case <-ticker.C:
|
||||
if ctx.Err() != nil {
|
||||
return nil
|
||||
}
|
||||
c.retryHeld(ctx, act)
|
||||
case delivery, ok := <-catchups:
|
||||
if !ok {
|
||||
if catchups != nil {
|
||||
return errors.New("the bus stopped delivering catch-up requests")
|
||||
}
|
||||
continue
|
||||
}
|
||||
act(ctx, c.wrap(KindCatchUp, delivery))
|
||||
case delivery, ok := <-upgrades:
|
||||
// A nil channel blocks for ever, so this case simply never fires when nothing is
|
||||
// listening for upgrades. Closed is different, and means the bus stopped.
|
||||
if !ok {
|
||||
if upgrades != nil {
|
||||
return errors.New("the bus stopped delivering upgrades")
|
||||
}
|
||||
continue
|
||||
}
|
||||
act(ctx, c.wrap(KindModuleMoved, delivery))
|
||||
case reason := <-closed:
|
||||
// Said rather than returned quietly. A controller whose bus connection dropped is a
|
||||
// mesh where nothing can be told anything, and the reason is the first thing anybody
|
||||
// will want.
|
||||
return fmt.Errorf("the bus connection closed: %v", reason)
|
||||
case delivery, ok := <-deliveries:
|
||||
if !ok {
|
||||
return errors.New("the bus stopped delivering")
|
||||
}
|
||||
kind, known := kindOfKey[delivery.RoutingKey]
|
||||
if !known {
|
||||
// Rejected without requeue: a message nothing understands will not be understood
|
||||
// on the next attempt either, and requeuing it would spin.
|
||||
_ = delivery.Reject(false)
|
||||
continue
|
||||
}
|
||||
act(ctx, c.wrap(kind, delivery))
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// kindOfKey is how this transport's addressing becomes what the mesh calls a message.
|
||||
var kindOfKey = map[string]string{
|
||||
KeyEnrol: KindEnrolment,
|
||||
KeyReport: KindReport,
|
||||
KeyAlive: KindHeartbeat,
|
||||
KeyBuilt: KindBuilt,
|
||||
KeyModuleUpgraded: KindModuleMoved,
|
||||
KeyCatchingUp: KindCatchUp,
|
||||
}
|
||||
|
||||
func (c *currentInbound) wrap(kind string, delivery amqp.Delivery) *currentControl {
|
||||
return ¤tControl{kind: kind, delivery: delivery, on: c}
|
||||
}
|
||||
|
||||
// retryHeld hands every message whose delay has passed back to the loop. Each handler holds it
|
||||
// again, settles it, or lets it go past the bound.
|
||||
func (c *currentInbound) retryHeld(ctx context.Context, act func(context.Context, Control)) {
|
||||
now := time.Now()
|
||||
due := make([]*currentControl, 0, len(c.held))
|
||||
for _, h := range c.held {
|
||||
if !h.due.After(now) {
|
||||
due = append(due, h.message)
|
||||
}
|
||||
}
|
||||
for _, m := range due {
|
||||
if ctx.Err() != nil {
|
||||
return
|
||||
}
|
||||
act(ctx, m)
|
||||
}
|
||||
}
|
||||
|
||||
// currentControl is one delivery from the bus the mesh has, as the controller reads it.
|
||||
type currentControl struct {
|
||||
kind string
|
||||
delivery amqp.Delivery
|
||||
on *currentInbound
|
||||
// about is what this message is about, as the handler named it; empty until it does.
|
||||
about string
|
||||
// first is when this message was first held for the store; zero while it has not been.
|
||||
first time.Time
|
||||
}
|
||||
|
||||
func (m *currentControl) Kind() string { return m.kind }
|
||||
func (m *currentControl) Body() []byte { return m.delivery.Body }
|
||||
func (m *currentControl) Redelivered() bool { return m.delivery.Redelivered }
|
||||
|
||||
func (m *currentControl) HeldFor() time.Duration {
|
||||
if m.first.IsZero() {
|
||||
return 0
|
||||
}
|
||||
return time.Since(m.first)
|
||||
}
|
||||
|
||||
// Answer publishes to the reply queue the request named.
|
||||
func (m *currentControl) Answer(ctx context.Context, body []byte) error {
|
||||
if m.delivery.ReplyTo == "" {
|
||||
return errors.New("that request named no reply queue, so nothing can be told the answer")
|
||||
}
|
||||
return m.on.channel.PublishWithContext(ctx, "", m.delivery.ReplyTo, false, false,
|
||||
amqp.Publishing{
|
||||
ContentType: "application/json",
|
||||
CorrelationId: m.delivery.CorrelationId,
|
||||
Body: body,
|
||||
})
|
||||
}
|
||||
|
||||
func (m *currentControl) Took() error {
|
||||
m.forget()
|
||||
return m.delivery.Ack(false)
|
||||
}
|
||||
|
||||
// Drop rejects without requeue: on this bus that is what "understood, and not worth another
|
||||
// attempt" is spelled as, and it is what feeds a dead-letter queue where one is configured.
|
||||
func (m *currentControl) Drop() error {
|
||||
m.forget()
|
||||
return m.delivery.Reject(false)
|
||||
}
|
||||
|
||||
// About names what this message is about, and lets go of whatever is held about the same thing:
|
||||
// the held one is the past, and acting on it after this one would undo this one. Acknowledged
|
||||
// rather than left to come back, because a held message nothing will act on is a place in the
|
||||
// prefetch nothing gets back.
|
||||
func (m *currentControl) About(what string) {
|
||||
m.about = what
|
||||
if what == "" {
|
||||
return
|
||||
}
|
||||
for tag, h := range m.on.held {
|
||||
if h.about != what || tag == m.delivery.DeliveryTag {
|
||||
continue
|
||||
}
|
||||
delete(m.on.held, tag)
|
||||
_ = h.message.delivery.Ack(false)
|
||||
}
|
||||
}
|
||||
|
||||
// Hold keeps the message unacknowledged and sets it aside to be handed back after the delay.
|
||||
//
|
||||
// Held no further than the prefetch leaves room: past that the bus would hand the loop nothing
|
||||
// new — enrolments included — until something held was let go. A message that cannot be held says
|
||||
// so, and the handler settles it its own way.
|
||||
func (m *currentControl) Hold(after time.Duration) error {
|
||||
if m.on.held == nil {
|
||||
m.on.held = map[uint64]*holding{}
|
||||
}
|
||||
if _, already := m.on.held[m.delivery.DeliveryTag]; !already {
|
||||
if len(m.on.held) >= Prefetch-PrefetchHeadroom {
|
||||
return fmt.Errorf("%d messages are already held for the store, and holding more "+
|
||||
"would stop the queue", len(m.on.held))
|
||||
}
|
||||
m.first = time.Now()
|
||||
}
|
||||
m.on.held[m.delivery.DeliveryTag] = &holding{
|
||||
message: m, due: time.Now().Add(after), about: m.about,
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (m *currentControl) forget() {
|
||||
if m.on != nil {
|
||||
delete(m.on.held, m.delivery.DeliveryTag)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,71 @@
|
||||
package link
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"io"
|
||||
"log"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
amqp "github.com/rabbitmq/amqp091-go"
|
||||
)
|
||||
|
||||
// The harness for the consume side on the bus the mesh runs on today.
|
||||
//
|
||||
// Messages arrive through the seam, so what these tests exercise is the controller's decision
|
||||
// about a message and this transport's way of keeping one — which is what the seam separated. A
|
||||
// fake acknowledger stands in for the bus, because what is asserted is how a message was settled
|
||||
// and that needs no server.
|
||||
|
||||
// settled is how the bus was told to settle one message.
|
||||
type settled struct{ acked, nacked, requeued, rejected bool }
|
||||
|
||||
func (a *settled) Ack(uint64, bool) error { a.acked = true; return nil }
|
||||
func (a *settled) Nack(_ uint64, _ bool, requeue bool) error {
|
||||
a.nacked, a.requeued = true, requeue
|
||||
return nil
|
||||
}
|
||||
func (a *settled) Reject(uint64, bool) error { a.rejected = true; return nil }
|
||||
|
||||
// unsettled is a message the controller has neither taken nor let go: it is held, and the bus will
|
||||
// hand it to whatever consumes next if the controller stops.
|
||||
func (a *settled) unsettled() bool { return !a.acked && !a.nacked && !a.rejected }
|
||||
|
||||
var tag uint64
|
||||
|
||||
func quiet() *log.Logger { return log.New(io.Discard, "", 0) }
|
||||
|
||||
// serving is a controller with nothing but a way of receiving, ready for a listener, a recorder,
|
||||
// an upgrader or a replayer to be set on it.
|
||||
func serving() (*Server, *currentInbound) {
|
||||
in := ¤tInbound{held: map[uint64]*holding{}}
|
||||
return &Server{inbound: in, bus: OverCurrent{}, log: quiet()}, in
|
||||
}
|
||||
|
||||
// sends is one message arriving over this transport, as the controller reads it.
|
||||
func (c *currentInbound) sends(t *testing.T, to *settled, kind string, v any) Control {
|
||||
t.Helper()
|
||||
body, err := json.Marshal(v)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
tag++
|
||||
return ¤tControl{kind: kind, on: c, delivery: amqp.Delivery{
|
||||
Acknowledger: to, Body: body, DeliveryTag: tag,
|
||||
}}
|
||||
}
|
||||
|
||||
// dueNow brings every held message forward, so a test need not wait out the backoff a real store
|
||||
// restart would be given (RedeliverAfter).
|
||||
func (c *currentInbound) dueNow() {
|
||||
for _, h := range c.held {
|
||||
h.due = time.Now().Add(-time.Second)
|
||||
}
|
||||
}
|
||||
|
||||
// retries hands every held message back to the controller, the way the ticker does.
|
||||
func (c *currentInbound) retries(ctx context.Context, s *Server) {
|
||||
c.dueNow()
|
||||
c.retryHeld(ctx, s.act)
|
||||
}
|
||||
@@ -1,94 +0,0 @@
|
||||
package link
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
// A harness for the controller's decisions about a message, with no bus behind it.
|
||||
//
|
||||
// What these tests exercise is the server's verdict — taken, dropped, held for the store, let go
|
||||
// once the store has been away too long — and the transport's part in that is only to keep a held
|
||||
// message and hand it back. A fake that does exactly that stands in for the bus, so what is
|
||||
// asserted is how a message was settled and the decision needs no server.
|
||||
|
||||
// settled is how the bus was told to settle one message.
|
||||
type settled struct{ acked, nacked, requeued, rejected bool }
|
||||
|
||||
// unsettled is a message the controller has neither taken nor let go: it is held, and the bus will
|
||||
// hand it to whatever consumes next if the controller stops.
|
||||
func (a *settled) unsettled() bool { return !a.acked && !a.nacked && !a.rejected }
|
||||
|
||||
// fakeInbound keeps the messages the controller held, the way a stream would.
|
||||
type fakeInbound struct{ held map[uint64]*fakeControl }
|
||||
|
||||
var tag uint64
|
||||
|
||||
// serving is a controller with nothing but a way of receiving, ready for a listener, a recorder,
|
||||
// an upgrader or a replayer to be set on it.
|
||||
func serving() (*Server, *fakeInbound) {
|
||||
in := &fakeInbound{held: map[uint64]*fakeControl{}}
|
||||
return &Server{inbound: in, log: quiet()}, in
|
||||
}
|
||||
|
||||
func (c *fakeInbound) Also(string) error { return nil }
|
||||
func (c *fakeInbound) Close() {}
|
||||
func (c *fakeInbound) Receive(context.Context, func(context.Context, Control)) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
// sends is one message arriving, as the controller reads it.
|
||||
func (c *fakeInbound) sends(t *testing.T, to *settled, kind string, v any) Control {
|
||||
t.Helper()
|
||||
body, err := json.Marshal(v)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
tag++
|
||||
return &fakeControl{kind: kind, body: body, tag: tag, to: to, on: c}
|
||||
}
|
||||
|
||||
// retries hands every held message back to the controller, the way a stream redelivers.
|
||||
func (c *fakeInbound) retries(ctx context.Context, s *Server) {
|
||||
for tag, m := range c.held {
|
||||
delete(c.held, tag)
|
||||
m.redelivered = true
|
||||
s.act(ctx, m)
|
||||
}
|
||||
}
|
||||
|
||||
type fakeControl struct {
|
||||
kind string
|
||||
body []byte
|
||||
tag uint64
|
||||
to *settled
|
||||
on *fakeInbound
|
||||
redelivered bool
|
||||
since time.Time
|
||||
}
|
||||
|
||||
func (m *fakeControl) Kind() string { return m.kind }
|
||||
func (m *fakeControl) Body() []byte { return m.body }
|
||||
func (m *fakeControl) Redelivered() bool { return m.redelivered }
|
||||
func (m *fakeControl) About(string) {}
|
||||
func (m *fakeControl) Answer(context.Context, []byte) error { return nil }
|
||||
func (m *fakeControl) Took() error { m.to.acked = true; return nil }
|
||||
func (m *fakeControl) Drop() error { m.to.rejected = true; return nil }
|
||||
|
||||
// HeldFor is how long the store has been waited on for this message — zero on a first delivery.
|
||||
func (m *fakeControl) HeldFor() time.Duration {
|
||||
if m.since.IsZero() {
|
||||
return 0
|
||||
}
|
||||
return time.Since(m.since)
|
||||
}
|
||||
|
||||
func (m *fakeControl) Hold(time.Duration) error {
|
||||
if m.since.IsZero() {
|
||||
m.since = time.Now()
|
||||
}
|
||||
m.on.held[m.tag] = m
|
||||
return nil
|
||||
}
|
||||
@@ -23,10 +23,6 @@ import (
|
||||
// it first could not take it, and a controller that restarts mid-window has nothing to lose.
|
||||
|
||||
// natsInbound consumes what nodes and modules say over NATS.
|
||||
// Prefetch is how many controls the controller holds unacknowledged at once: the store window
|
||||
// (ADR 0083) is the server's, and this is what it may hand this process ahead of its acting.
|
||||
const Prefetch = 64
|
||||
|
||||
type natsInbound struct {
|
||||
js *broker.JetStream
|
||||
// follows is the kinds asked for beyond what nodes say (Also). The events those are are the
|
||||
@@ -226,6 +222,14 @@ func (m *natsControl) HeldFor() time.Duration {
|
||||
return time.Since(first)
|
||||
}
|
||||
|
||||
// About is nothing here, and that is the point.
|
||||
//
|
||||
// Setting a held message aside when a newer one about the same thing arrives is what a controller
|
||||
// holding deliveries in memory can do. A naked message belongs to the server and comes back
|
||||
// whatever happened meanwhile, so the question "is this the past?" is answered by what the message
|
||||
// says instead — the digest of the declaration a report is about (window.go, design 25 §3).
|
||||
func (m *natsControl) About(string) {}
|
||||
|
||||
// Answer publishes to the reply subject the request carries **in its payload**.
|
||||
//
|
||||
// Not `Respond`, and not the message's reply field: a message a JetStream consumer delivers has had
|
||||
|
||||
@@ -4,8 +4,6 @@ import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"io"
|
||||
"log"
|
||||
"os"
|
||||
"sync"
|
||||
"testing"
|
||||
@@ -122,8 +120,6 @@ func eventually(t *testing.T, what string, is func() bool) {
|
||||
|
||||
// A report published by a node reaches the controller, is recorded, and is acknowledged — so the
|
||||
// stream does not hold it. A work queue is the check: what is acknowledged leaves it.
|
||||
func quiet() *log.Logger { return log.New(io.Discard, "", 0) }
|
||||
|
||||
func TestNatsAReportIsHeardAndLeavesTheStream(t *testing.T) {
|
||||
js := aBus(t)
|
||||
store := &counted{}
|
||||
|
||||
@@ -46,6 +46,24 @@ func TestAReportTheStoreCouldNotTakeIsHeldAndOneItRefusedIsNot(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
// A newer report from the same node supersedes one of its reports still held: recorded after the
|
||||
// newer, the older would overwrite what the node is doing now.
|
||||
func TestANewerReportSupersedesAHeldOneFromTheSameNode(t *testing.T) {
|
||||
s, in := serving()
|
||||
s.listener = heardWith{err: errors.Join(ErrTryAgain, errors.New("starting up"))}
|
||||
older, newer, other := &settled{}, &settled{}, &settled{}
|
||||
s.act(context.Background(), in.sends(t, older, KindReport, aReport("anchor", "d1")))
|
||||
s.act(context.Background(), in.sends(t, other, KindReport, aReport("laptop", "d7")))
|
||||
s.act(context.Background(), in.sends(t, newer, KindReport, aReport("anchor", "d2")))
|
||||
if !older.acked {
|
||||
t.Fatalf("the older report was not set aside by the newer: %+v", older)
|
||||
}
|
||||
if !newer.unsettled() || !other.unsettled() || len(in.held) != 2 {
|
||||
t.Fatalf("the newer report and another node's were not both held: newer %+v other %+v, %d held",
|
||||
newer, other, len(in.held))
|
||||
}
|
||||
}
|
||||
|
||||
// A store that has not come back within the bound is not restarting: the report is let go, loudly,
|
||||
// rather than held for ever.
|
||||
func TestAReportIsLetGoOnceTheStoreHasBeenGoneTooLong(t *testing.T) {
|
||||
|
||||
+84
-4
@@ -11,6 +11,10 @@ import (
|
||||
"log"
|
||||
"os"
|
||||
"time"
|
||||
|
||||
amqp "github.com/rabbitmq/amqp091-go"
|
||||
|
||||
"github.com/novox/mesh-controller/internal/envfile"
|
||||
)
|
||||
|
||||
// AMQPVar is the controller's own connection to the bus the mesh runs on today.
|
||||
@@ -67,6 +71,8 @@ type Upgrader interface {
|
||||
type Server struct {
|
||||
inbound Inbound
|
||||
bus Bus
|
||||
conn *amqp.Connection
|
||||
channel *amqp.Channel
|
||||
js *broker.JetStream
|
||||
|
||||
enroller Enroller
|
||||
@@ -121,18 +127,85 @@ func (s *Server) Answers(r Replayer) error {
|
||||
// On the port MESH_BROKER_AMQP_PORT names when the node's settings moved the broker (novox/hq
|
||||
// 04-ISSUES/102) — the URL is genesis's, sealed, and its port is the one thing in it the node may
|
||||
// have moved since.
|
||||
// Connect is ConnectNats: the mesh has one bus (novox/hq ADR 0131, design 28 task 5.5). Kept as
|
||||
// the name callers know; the transport it opened before is gone with the bus it spoke to.
|
||||
func Connect(js *broker.JetStream, enroller Enroller, listener Listener) *Server {
|
||||
return ConnectNats(js, enroller, listener)
|
||||
func Connect(enroller Enroller, listener Listener) (*Server, error) {
|
||||
url, err := envfile.Placed(AMQPVar)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if url == "" {
|
||||
return nil, fmt.Errorf(
|
||||
"this control plane has no %s, so it cannot reach its broker. Nodes talk to it over "+
|
||||
"the broker and nowhere else, so without this it can hold records and answer "+
|
||||
"nothing", AMQPVar)
|
||||
}
|
||||
|
||||
conn, err := amqp.Dial(url)
|
||||
if err != nil {
|
||||
// Not quoted back: the URL carries the controller's own bus password.
|
||||
return nil, fmt.Errorf("cannot reach the broker named in %s: %w", AMQPVar, err)
|
||||
}
|
||||
channel, err := conn.Channel()
|
||||
if err != nil {
|
||||
conn.Close()
|
||||
return nil, err
|
||||
}
|
||||
|
||||
// Declared here rather than assumed. The controller is the only thing that may create them — a
|
||||
// node's account can write to this exchange and read its own queue, and configure nothing
|
||||
// else, so a node arriving before the controller has ever run finds nothing and says so,
|
||||
// rather than quietly creating a topology nobody designed.
|
||||
if err := channel.ExchangeDeclare(Exchange, "direct", true, false, false, false, nil); err != nil {
|
||||
conn.Close()
|
||||
return nil, fmt.Errorf("cannot declare the %s exchange: %w", Exchange, err)
|
||||
}
|
||||
// The events exchange too. The controller is not the only publisher on it — modules announce
|
||||
// onto it with their own accounts — but it is the only thing permitted to create it, for the
|
||||
// same reason it is the only thing permitted to create the direct one.
|
||||
if err := channel.ExchangeDeclare(EventsExchange, "topic", true, false, false, false, nil); err != nil {
|
||||
conn.Close()
|
||||
return nil, fmt.Errorf("cannot declare the %s exchange: %w", EventsExchange, err)
|
||||
}
|
||||
if _, err := channel.QueueDeclare(ControlQueue, true, false, false, false, nil); err != nil {
|
||||
conn.Close()
|
||||
return nil, fmt.Errorf("cannot declare the %s queue: %w", ControlQueue, err)
|
||||
}
|
||||
// Every key a node may publish. Binding one and forgetting another is a message the bus
|
||||
// accepts, finds no queue for, and drops — the publisher sees success and the consumer sees
|
||||
// nothing. That is exactly what happened to reports: `report` was left unbound while `enrol`
|
||||
// worked, so nodes announced what they had applied into a void for an afternoon.
|
||||
for _, key := range []string{KeyEnrol, KeyReport, KeyAlive, KeyBuilt} {
|
||||
if err := channel.QueueBind(ControlQueue, key, Exchange, false, nil); err != nil {
|
||||
conn.Close()
|
||||
return nil, fmt.Errorf("cannot bind %s to %s/%s: %w", ControlQueue, Exchange, key, err)
|
||||
}
|
||||
}
|
||||
|
||||
return &Server{
|
||||
inbound: Current(conn, channel),
|
||||
bus: OverCurrent{Channel: channel},
|
||||
conn: conn,
|
||||
channel: channel,
|
||||
enroller: enroller,
|
||||
listener: listener,
|
||||
log: newLog(),
|
||||
}, nil
|
||||
}
|
||||
|
||||
func newLog() *log.Logger { return log.New(os.Stdout, "", log.LstdFlags) }
|
||||
|
||||
// Channel is the controller's channel, for the command line's own publishing.
|
||||
func (s *Server) Channel() *amqp.Channel { return s.channel }
|
||||
|
||||
func (s *Server) Close() {
|
||||
if s.inbound != nil {
|
||||
s.inbound.Close()
|
||||
}
|
||||
if s.channel != nil {
|
||||
_ = s.channel.Close()
|
||||
}
|
||||
if s.conn != nil {
|
||||
_ = s.conn.Close()
|
||||
}
|
||||
if s.js != nil {
|
||||
s.js.Close()
|
||||
}
|
||||
@@ -276,6 +349,7 @@ func (s *Server) reported(ctx context.Context, m Control) {
|
||||
_ = m.Drop()
|
||||
return
|
||||
}
|
||||
m.About("report " + report.Node)
|
||||
what := fmt.Sprintf("%s's report of declaration %s", report.Node, report.Declared)
|
||||
|
||||
if s.listener != nil {
|
||||
@@ -408,6 +482,9 @@ func (s *Server) wasBuilt(ctx context.Context, m Control) {
|
||||
_ = m.Drop()
|
||||
return
|
||||
}
|
||||
// Each build result its own subject: none supersedes another, and recording one twice is
|
||||
// harmless — the build is kept by its id.
|
||||
m.About("build " + digest(m.Body()))
|
||||
err := s.recorder.Built(ctx, result)
|
||||
switch s.decide(ctx, m, fmt.Sprintf("a build result from %s", result.On), "", "", err) {
|
||||
case Hold:
|
||||
@@ -439,6 +516,7 @@ func (s *Server) wasBuilt(ctx context.Context, m Control) {
|
||||
// the catalogue next restarts (issue 083). One request stands for all: a newer one supersedes one
|
||||
// still held.
|
||||
func (s *Server) catchingUp(ctx context.Context, m Control) {
|
||||
m.About("catch-up")
|
||||
if s.replayer == nil {
|
||||
s.log.Printf("a catalogue asked to catch up and this control plane has nothing to replay")
|
||||
_ = m.Took()
|
||||
@@ -491,6 +569,7 @@ func (s *Server) moved(ctx context.Context, m Control) {
|
||||
_ = m.Took()
|
||||
return
|
||||
}
|
||||
m.About("upgrade " + u.Module)
|
||||
if u.Module == "" {
|
||||
s.log.Printf("an upgrade announcement named no module; ignored")
|
||||
_ = m.Took()
|
||||
@@ -541,6 +620,7 @@ func (s *Server) sourceMoved(ctx context.Context, m Control) {
|
||||
_ = m.Took()
|
||||
return
|
||||
}
|
||||
m.About("merge " + moved.Owner + "/" + moved.Repo + " into " + moved.Base)
|
||||
if moved.Commit == "" || moved.Repo == "" {
|
||||
s.log.Printf("a merge announcement named no repository or no commit; ignored")
|
||||
_ = m.Took()
|
||||
|
||||
@@ -3,6 +3,7 @@ package link
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"testing"
|
||||
|
||||
"github.com/jackc/pgx/v5/pgconn"
|
||||
@@ -17,8 +18,7 @@ func (r recordsWith) Built(context.Context, BuildResult) error { return r.err }
|
||||
|
||||
type upgradesWith struct{ err error }
|
||||
|
||||
func (u upgradesWith) Upgraded(context.Context, Upgraded) error { return u.err }
|
||||
func (u upgradesWith) SourceMoved(context.Context, SourceMoved) error { return u.err }
|
||||
func (u upgradesWith) Upgraded(context.Context, Upgraded) error { return u.err }
|
||||
|
||||
type replaysWith struct{ err error }
|
||||
|
||||
@@ -106,3 +106,51 @@ func TestAMessageHandledDuringShutdownIsLeftForTheBus(t *testing.T) {
|
||||
t.Fatalf("a build result handled during shutdown was settled, and so lost: %+v", *to)
|
||||
}
|
||||
}
|
||||
|
||||
// Two identical build results: the newer sets the older aside rather than leaving it unsettled
|
||||
// for ever, holding a place in the prefetch.
|
||||
func TestAnIdenticalBuildResultSetsTheHeldOneAside(t *testing.T) {
|
||||
s, in := serving()
|
||||
s.recorder = recordsWith{err: restarting}
|
||||
built := BuildResult{On: "anchor", Repository: "/r", Commit: "abc"}
|
||||
first, second := &settled{}, &settled{}
|
||||
s.act(context.Background(), in.sends(t, first, KindBuilt, built))
|
||||
s.act(context.Background(), in.sends(t, second, KindBuilt, built))
|
||||
if !first.acked || !second.unsettled() || len(in.held) != 1 {
|
||||
t.Fatalf("an identical build result did not set the held one aside: first %+v second %+v, %d held",
|
||||
*first, *second, len(in.held))
|
||||
}
|
||||
}
|
||||
|
||||
// What is held stops short of the prefetch, so the loop always has room to answer an enrolment.
|
||||
func TestWhatIsHeldLeavesRoomInThePrefetch(t *testing.T) {
|
||||
s, in := serving()
|
||||
s.recorder = recordsWith{err: restarting}
|
||||
var last *settled
|
||||
for i := 0; i < Prefetch; i++ {
|
||||
last = &settled{}
|
||||
s.act(context.Background(), in.sends(t, last, KindBuilt, BuildResult{On: "anchor", Commit: fmt.Sprint(i)}))
|
||||
}
|
||||
if len(in.held) != Prefetch-PrefetchHeadroom {
|
||||
t.Fatalf("%d messages were held; the ceiling is %d", len(in.held), Prefetch-PrefetchHeadroom)
|
||||
}
|
||||
if last.unsettled() {
|
||||
t.Fatalf("a message past the ceiling was held: %+v", *last)
|
||||
}
|
||||
}
|
||||
|
||||
// An upgrade handled during shutdown is left for the bus too — the upgrader's error is the
|
||||
// cancelled context, which is no answer about the announcement.
|
||||
func TestAnUpgradeHandledDuringShutdownIsLeftForTheBus(t *testing.T) {
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
cancel()
|
||||
s, in := serving()
|
||||
s.upgrader = upgradesWith{err: context.Canceled}
|
||||
to := &settled{}
|
||||
s.act(ctx, in.sends(t, to, KindModuleMoved, Upgraded{Module: "gitea", Commit: "abcdef0123"}))
|
||||
if !to.unsettled() {
|
||||
t.Fatalf("an upgrade was settled during shutdown, and so lost: %+v", *to)
|
||||
}
|
||||
}
|
||||
|
||||
func (u upgradesWith) SourceMoved(context.Context, SourceMoved) error { return nil }
|
||||
|
||||
Reference in New Issue
Block a user