Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
c585158836 |
@@ -37,7 +37,7 @@ func askCommand(ctx context.Context, args []string) error {
|
|||||||
arguments = json.RawMessage(positionals[2])
|
arguments = json.RawMessage(positionals[2])
|
||||||
}
|
}
|
||||||
|
|
||||||
server, err := connectLink(ctx, nil, nil, nil)
|
server, err := link.Connect(nil, nil)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -368,7 +368,7 @@ func buildOne(ctx context.Context, source buildSource, path, ref string, wait ti
|
|||||||
}
|
}
|
||||||
defer ident.Close()
|
defer ident.Close()
|
||||||
|
|
||||||
server, err := connectLink(ctx, nil, nil, nil)
|
server, err := link.Connect(nil, nil)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
@@ -472,7 +472,7 @@ func buildAndShow(ctx context.Context, source buildSource, path, ref string, wai
|
|||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
defer ident.Close()
|
defer ident.Close()
|
||||||
server, err := connectLink(ctx, nil, nil, nil)
|
server, err := link.Connect(nil, nil)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -164,6 +164,8 @@ func usage() {
|
|||||||
upgrade <name> record ...record that they are behind, and send nothing
|
upgrade <name> record ...record that they are behind, and send nothing
|
||||||
status [--json] what is wrong, what is quiet, and what is out of date
|
status [--json] what is wrong, what is quiet, and what is out of date
|
||||||
seats [--json] every seat this mesh defines, what it delivers, and who holds it
|
seats [--json] every seat this mesh defines, what it delivers, and who holds it
|
||||||
|
seat rename <from> <to> rename a seat; its former name still resolves (ADR 0122)
|
||||||
|
seat <name> --to <node>/<module> hand a seat to that assignment as one act; never empty in between (ADR 0131)
|
||||||
board [--listen ADDR] the same three questions, as a page that holds nothing
|
board [--listen ADDR] the same three questions, as a page that holds nothing
|
||||||
api --issuer URL [--listen A] assign and unassign over http, for a surface that is not here
|
api --issuer URL [--listen A] assign and unassign over http, for a surface that is not here
|
||||||
assign <node> <module> put a module on a node
|
assign <node> <module> put a module on a node
|
||||||
|
|||||||
@@ -614,15 +614,6 @@ func issueOnTheNewBus(ctx context.Context, inv *inventory.Inventory, m catalogue
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
return issueWith(ctx, inv, m, node, busAddress, known, reachable, user, password)
|
|
||||||
}
|
|
||||||
|
|
||||||
// issueWith is the delivery half: the minted password sealed to the machine as the module's broker
|
|
||||||
// secret, and the module's consumer created where the bus can be reached. Split from the minting
|
|
||||||
// so the move can issue every module against a bus whose address it worked out itself
|
|
||||||
// (`rollout mint`, design 28 task 5.2) rather than the one in this process's environment.
|
|
||||||
func issueWith(ctx context.Context, inv *inventory.Inventory, m catalogue.Manifest,
|
|
||||||
node, busAddress string, known broker.Broker, reachable, user, password string) error {
|
|
||||||
held, err := json.Marshal(struct {
|
held, err := json.Marshal(struct {
|
||||||
URL string `json:"url"`
|
URL string `json:"url"`
|
||||||
Fingerprint string `json:"fingerprint,omitempty"`
|
Fingerprint string `json:"fingerprint,omitempty"`
|
||||||
@@ -648,19 +639,14 @@ func issueWith(ctx context.Context, inv *inventory.Inventory, m catalogue.Manife
|
|||||||
Kind: broker.KindModule, Node: node, Module: m.Module,
|
Kind: broker.KindModule, Node: node, Module: m.Module,
|
||||||
Emits: m.Emits, Consumes: m.Consumes, Serves: m.Tools,
|
Emits: m.Emits, Consumes: m.Consumes, Serves: m.Tools,
|
||||||
}); needed {
|
}); needed {
|
||||||
if busAddress == "" {
|
js, err := broker.Dial(busAddress)
|
||||||
fmt.Printf(" %s consumes; its consumer is created when the bus is reachable (`push`, then "+
|
if err != nil {
|
||||||
"`rollout mint` again is harmless)\n", m.Module)
|
return fmt.Errorf("the credential is minted and the mesh cannot reach the bus to create "+
|
||||||
} else {
|
"how %s hears what it consumes: %w", m.Module, err)
|
||||||
js, err := broker.Dial(busAddress)
|
}
|
||||||
if err != nil {
|
defer js.Close()
|
||||||
return fmt.Errorf("the credential is minted and the mesh cannot reach the bus to create "+
|
if err := js.EnsureConsumer(consumer); err != nil {
|
||||||
"how %s hears what it consumes: %w", m.Module, err)
|
return err
|
||||||
}
|
|
||||||
defer js.Close()
|
|
||||||
if err := js.EnsureConsumer(consumer); err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -537,15 +537,6 @@ func onTheNetwork(ctx context.Context, inv *inventory.Inventory,
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
// **With the seat holders on record**, or a machine running the next holder of a seat beside
|
|
||||||
// the current one resolves as two holders, is refused, and drops out of the map — taking the
|
|
||||||
// address every other machine composes for what it offers (novox/hq ADR 0131). Found live:
|
|
||||||
// the control node vanished from the private network the moment the new bus was assigned
|
|
||||||
// beside the old one.
|
|
||||||
holdings, err := inv.Holdings(ctx)
|
|
||||||
if err != nil {
|
|
||||||
return nil, err
|
|
||||||
}
|
|
||||||
var out []inventory.Overlay
|
var out []inventory.Overlay
|
||||||
for _, p := range places {
|
for _, p := range places {
|
||||||
if p.Address == "" {
|
if p.Address == "" {
|
||||||
@@ -558,7 +549,7 @@ func onTheNetwork(ctx context.Context, inv *inventory.Inventory,
|
|||||||
caps, _ := inv.ProfileOf(ctx, p.Name)
|
caps, _ := inv.ProfileOf(ctx, p.Name)
|
||||||
got, err := catalogue.Resolve(shelf, assigned,
|
got, err := catalogue.Resolve(shelf, assigned,
|
||||||
catalogue.Node{Name: p.Name, Site: p.Site, Capabilities: caps},
|
catalogue.Node{Name: p.Name, Site: p.Site, Capabilities: caps},
|
||||||
catalogue.World{Unchecked: true, Holdings: holdings})
|
catalogue.World{Unchecked: true})
|
||||||
if err != nil {
|
if err != nil {
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -636,13 +636,8 @@ func renderingFor(ctx context.Context, open *stores, node string,
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
return catalogue.Rendering{}, inventory.Node{}, err
|
return catalogue.Rendering{}, inventory.Node{}, err
|
||||||
}
|
}
|
||||||
memberships, err := inv.BusMemberships(ctx)
|
|
||||||
if err != nil {
|
|
||||||
return catalogue.Rendering{}, inventory.Node{}, err
|
|
||||||
}
|
|
||||||
return catalogue.Rendering{
|
return catalogue.Rendering{
|
||||||
BusMembership: memberships[node],
|
Settings: settings, Generators: gens, Grants: grants, Needed: needed, Ports: ports,
|
||||||
Settings: settings, Generators: gens, Grants: grants, Needed: needed, Ports: ports,
|
|
||||||
Certificate: certificate, Authority: authority, Mesh: private, Names: names,
|
Certificate: certificate, Authority: authority, Mesh: private, Names: names,
|
||||||
Machines: machines,
|
Machines: machines,
|
||||||
Suffix: overlay.Suffix(), MeshRange: meshRange, Accounts: accounts, Foundation: foundation,
|
Suffix: overlay.Suffix(), MeshRange: meshRange, Accounts: accounts, Foundation: foundation,
|
||||||
|
|||||||
+14
-36
@@ -38,33 +38,6 @@ func reportUnhostable(node string, plan catalogue.Resolution) {
|
|||||||
// nothing in it was wrong, and no one edit was the one that should have been a new file.
|
// nothing in it was wrong, and no one edit was the one that should have been a new file.
|
||||||
|
|
||||||
// serve is the control plane running: one connection to the broker, one queue, one consumer.
|
// serve is the control plane running: one connection to the broker, one queue, one consumer.
|
||||||
// connectLink opens the controller's link over whichever bus this process is on (design 25: one
|
|
||||||
// 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, 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
|
|
||||||
}
|
|
||||||
}
|
|
||||||
js, err := broker.Dial(busAddress)
|
|
||||||
if err != nil {
|
|
||||||
return nil, fmt.Errorf("the mesh is on the bus at %s and this control plane cannot reach it: %w",
|
|
||||||
broker.BareAddress(busAddress), err)
|
|
||||||
}
|
|
||||||
return link.ConnectNats(js, enroller, listener), nil
|
|
||||||
}
|
|
||||||
|
|
||||||
func serve(ctx context.Context) error {
|
func serve(ctx context.Context) error {
|
||||||
open, err := openStores(ctx)
|
open, err := openStores(ctx)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@@ -117,7 +90,7 @@ func serve(ctx context.Context) error {
|
|||||||
|
|
||||||
work := link.Enrolment{Inventory: inv, Identity: ident, Management: management, Broker: known,
|
work := link.Enrolment{Inventory: inv, Identity: ident, Management: management, Broker: known,
|
||||||
OnNATS: onNATS}
|
OnNATS: onNATS}
|
||||||
server, err := connectLink(ctx, inv, work, work)
|
server, err := link.Connect(work, work)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
@@ -127,6 +100,11 @@ func serve(ctx context.Context) error {
|
|||||||
// somebody deleted, a mesh raised from a restored backup, or a bus whose data directory was
|
// somebody deleted, a mesh raised from a restored backup, or a bus whose data directory was
|
||||||
// replaced all have records and no objects — and a node whose consumer is missing hears nothing
|
// replaced all have records and no objects — and a node whose consumer is missing hears nothing
|
||||||
// while everything else about it looks correct.
|
// while everything else about it looks correct.
|
||||||
|
if onNATS {
|
||||||
|
if err := raiseTheBus(ctx, inv, busAddress); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
}
|
||||||
// And build results nobody was waiting for. A build triggered any other way than `build`
|
// And build results nobody was waiting for. A build triggered any other way than `build`
|
||||||
// would otherwise be reported into the void, which is the same as not reporting it.
|
// would otherwise be reported into the void, which is the same as not reporting it.
|
||||||
server.Records(builds{inv})
|
server.Records(builds{inv})
|
||||||
@@ -180,13 +158,13 @@ func declare(ctx context.Context, args []string) error {
|
|||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
||||||
server, err := connectLink(ctx, nil, nil, nil)
|
server, err := link.Connect(nil, nil)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
defer server.Close()
|
defer server.Close()
|
||||||
|
|
||||||
if err := link.Declare(ctx, server.Bus(), ident, node, raw, 15*time.Second); err != nil {
|
if err := link.Declare(ctx, link.OverCurrent{Channel: server.Channel()}, ident, node, raw, 15*time.Second); err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
fmt.Printf("sent %s a signed declaration (%d bytes)\n", node, len(raw))
|
fmt.Printf("sent %s a signed declaration (%d bytes)\n", node, len(raw))
|
||||||
@@ -293,7 +271,7 @@ func pushCommand(ctx context.Context, args []string) error {
|
|||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
||||||
server, err := connectLink(ctx, nil, nil, nil)
|
server, err := link.Connect(nil, nil)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
@@ -354,7 +332,7 @@ func pushCommand(ctx context.Context, args []string) error {
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
if err := link.Declare(ctx, server.Bus(), ident, s.node, body, 15*time.Second); err != nil {
|
if err := link.Declare(ctx, link.OverCurrent{Channel: server.Channel()}, ident, s.node, body, 15*time.Second); err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
// After it is away, not before. A digest recorded for something that failed to send would
|
// After it is away, not before. A digest recorded for something that failed to send would
|
||||||
@@ -437,7 +415,7 @@ func pushCommand(ctx context.Context, args []string) error {
|
|||||||
return declarationWith(held, open, node, plan, settings, gens, Allocating)
|
return declarationWith(held, open, node, plan, settings, gens, Allocating)
|
||||||
},
|
},
|
||||||
func(s readyNode, body []byte) error {
|
func(s readyNode, body []byte) error {
|
||||||
if err := link.Declare(ctx, server.Bus(), ident, s.node, body,
|
if err := link.Declare(ctx, link.OverCurrent{Channel: server.Channel()}, ident, s.node, body,
|
||||||
15*time.Second); err != nil {
|
15*time.Second); err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
@@ -650,7 +628,7 @@ func sendTo(ctx context.Context, open *stores, names []string) error {
|
|||||||
len(refusals), strings.Join(refusals, "\n\n"))
|
len(refusals), strings.Join(refusals, "\n\n"))
|
||||||
}
|
}
|
||||||
|
|
||||||
server, err := connectLink(ctx, nil, nil, nil)
|
server, err := link.Connect(nil, nil)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
@@ -661,7 +639,7 @@ func sendTo(ctx context.Context, open *stores, names []string) error {
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
if err := link.Declare(ctx, server.Bus(), ident, s.node, body, 15*time.Second); err != nil {
|
if err := link.Declare(ctx, link.OverCurrent{Channel: server.Channel()}, ident, s.node, body, 15*time.Second); err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
record, err := inv.NodeByName(ctx, s.node)
|
record, err := inv.NodeByName(ctx, s.node)
|
||||||
@@ -726,7 +704,7 @@ func raiseTheBus(ctx context.Context, inv *inventory.Inventory, address string)
|
|||||||
js, err := broker.Dial(address)
|
js, err := broker.Dial(address)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return fmt.Errorf("the mesh is on the bus at %s and this control plane cannot reach it: %w",
|
return fmt.Errorf("the mesh is on the bus at %s and this control plane cannot reach it: %w",
|
||||||
broker.BareAddress(address), err)
|
address, err)
|
||||||
}
|
}
|
||||||
defer js.Close()
|
defer js.Close()
|
||||||
|
|
||||||
|
|||||||
@@ -2,7 +2,6 @@ package main
|
|||||||
|
|
||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
"encoding/json"
|
|
||||||
"errors"
|
"errors"
|
||||||
"fmt"
|
"fmt"
|
||||||
"strings"
|
"strings"
|
||||||
@@ -13,7 +12,6 @@ import (
|
|||||||
"github.com/novox/mesh-controller/internal/broker"
|
"github.com/novox/mesh-controller/internal/broker"
|
||||||
"github.com/novox/mesh-controller/internal/catalogue"
|
"github.com/novox/mesh-controller/internal/catalogue"
|
||||||
"github.com/novox/mesh-controller/internal/inventory"
|
"github.com/novox/mesh-controller/internal/inventory"
|
||||||
"github.com/novox/mesh-controller/internal/secrets"
|
|
||||||
)
|
)
|
||||||
|
|
||||||
// Moving the mesh's own traffic to the bus being built (novox/hq ADR 0116 step 5).
|
// Moving the mesh's own traffic to the bus being built (novox/hq ADR 0116 step 5).
|
||||||
@@ -33,19 +31,12 @@ import (
|
|||||||
// ability to change things, not the services its modules are serving — measured on 2026-09-27, when
|
// ability to change things, not the services its modules are serving — measured on 2026-09-27, when
|
||||||
// a seat emptied mid-change and the control plane looped for two hours while every service stayed up.
|
// a seat emptied mid-change and the control plane looped for two hours while every service stayed up.
|
||||||
|
|
||||||
const rolloutUsage = "rollout check | rollout mint [--again] | rollout --confirm"
|
const rolloutUsage = "rollout check | rollout --confirm"
|
||||||
|
|
||||||
func rolloutCommand(ctx context.Context, args []string) error {
|
func rolloutCommand(ctx context.Context, args []string) error {
|
||||||
switch {
|
switch {
|
||||||
case len(args) == 1 && args[0] == "check":
|
case len(args) == 1 && args[0] == "check":
|
||||||
return rolloutCheck(ctx)
|
return rolloutCheck(ctx)
|
||||||
case len(args) == 1 && args[0] == "mint":
|
|
||||||
return rolloutMint(ctx, false)
|
|
||||||
case len(args) == 2 && args[0] == "mint" && args[1] == "--again":
|
|
||||||
// Every credential minted afresh, whether or not one exists — for a mint that was wrong
|
|
||||||
// before anything was pushed. Afterwards nothing that received the old one still works,
|
|
||||||
// which is fine exactly when nothing received it.
|
|
||||||
return rolloutMint(ctx, true)
|
|
||||||
case len(args) == 1 && args[0] == "--confirm":
|
case len(args) == 1 && args[0] == "--confirm":
|
||||||
return errors.New(
|
return errors.New(
|
||||||
"the rollout itself is not built yet: `rollout check` answers whether it could run, and " +
|
"the rollout itself is not built yet: `rollout check` answers whether it could run, and " +
|
||||||
@@ -199,191 +190,3 @@ func wasSentTheUserList(ctx context.Context, inv *inventory.Inventory, node stri
|
|||||||
// notReadyOf is the readiness reasoning, named here so a test can reach it without the command's
|
// notReadyOf is the readiness reasoning, named here so a test can reach it without the command's
|
||||||
// printing. The reasoning itself is the broker package's, where it is pure.
|
// printing. The reasoning itself is the broker package's, where it is pure.
|
||||||
func notReadyOf(state broker.Readiness) []string { return broker.NotReady(state) }
|
func notReadyOf(state broker.Readiness) []string { return broker.NotReady(state) }
|
||||||
|
|
||||||
// rolloutMint gives every principal the new bus will have a credential it does not yet have, and
|
|
||||||
// puts each where its owner reads it (novox/hq design 28, task 5.2): a machine's as a membership
|
|
||||||
// sealed into its declaration, a module's as its broker secret, the control plane's own as its
|
|
||||||
// `bus` secret. Idempotent: what already has a hash is left alone, so running it again is harmless.
|
|
||||||
//
|
|
||||||
// **Before anything moves, and it is what makes moving possible.** A machine moved without a
|
|
||||||
// credential cannot come back, and afterwards there is no bus to tell it anything over — which is
|
|
||||||
// why `rollout check` refuses until this has run. The bus's address is worked out here, from where
|
|
||||||
// the module that provides it is assigned, rather than read from this process's environment: this
|
|
||||||
// process is still on the old bus when this runs, and must be.
|
|
||||||
func rolloutMint(ctx context.Context, again bool) error {
|
|
||||||
open, err := openStores(ctx)
|
|
||||||
if err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
defer open.Close()
|
|
||||||
inv := open.inventory
|
|
||||||
|
|
||||||
known, err := broker.FromEnvironment()
|
|
||||||
if err != nil {
|
|
||||||
return fmt.Errorf("the bus's certificate is not known to this process, and every membership "+
|
|
||||||
"must carry its fingerprint: %w", err)
|
|
||||||
}
|
|
||||||
shelf, err := inv.Catalogue(ctx)
|
|
||||||
if err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
entries, err := inv.Catalogued(ctx)
|
|
||||||
if err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
var busNode, controllerNode string
|
|
||||||
for _, e := range entries {
|
|
||||||
switch {
|
|
||||||
case e.Manifest.ClaimsSeat("mesh-broker") && providesBus(e.Manifest) && len(e.On) > 0:
|
|
||||||
busNode = e.On[0]
|
|
||||||
case e.Manifest.Module == "mesh-controller" && len(e.On) > 0:
|
|
||||||
controllerNode = e.On[0]
|
|
||||||
}
|
|
||||||
}
|
|
||||||
if busNode == "" {
|
|
||||||
return errors.New("no assigned module provides mesh-bus and claims mesh-broker, so there is no " +
|
|
||||||
"bus to mint credentials for — register and assign it first")
|
|
||||||
}
|
|
||||||
onNetwork, err := whereEveryoneIs(ctx, inv, shelf)
|
|
||||||
if err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
busHost := onNetwork[busNode]
|
|
||||||
if busHost == "" {
|
|
||||||
// **The hub is not in that map.** The machine that took over the tunnel is where the current
|
|
||||||
// bus already answers, and every machine dials it at the address the mesh handed them — so
|
|
||||||
// when the new bus runs on the same machine, that address is the one to tell them, with the
|
|
||||||
// new port. Found live: the control node is the hub, and the map lists the machines placed
|
|
||||||
// around it.
|
|
||||||
// The host alone: no scheme (BareAddress adds one where none was, which is the wrong
|
|
||||||
// direction here — every URL built below adds its own) and no port.
|
|
||||||
_, _, host := broker.CredentialIn(known.Address)
|
|
||||||
if host == "" {
|
|
||||||
host = known.Address
|
|
||||||
}
|
|
||||||
if _, after, hasScheme := strings.Cut(host, "://"); hasScheme {
|
|
||||||
host = after
|
|
||||||
}
|
|
||||||
host = strings.TrimSpace(host)
|
|
||||||
if i := strings.LastIndex(host, ":"); i > 0 && !strings.Contains(host[i:], "]") {
|
|
||||||
host = host[:i]
|
|
||||||
}
|
|
||||||
if host == "" {
|
|
||||||
return fmt.Errorf("%s runs the new bus and has no address on the private network, and the "+
|
|
||||||
"current bus's address is unknown too, so no machine could be told where it is", busNode)
|
|
||||||
}
|
|
||||||
busHost = host
|
|
||||||
}
|
|
||||||
busAddress := busHost + ":4222"
|
|
||||||
|
|
||||||
records, err := inv.BusRecords(ctx)
|
|
||||||
if err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
users, err := broker.Users(records)
|
|
||||||
if err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
kept, err := inv.BusUsers(ctx)
|
|
||||||
if err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
hashes := make(map[string]string, len(kept))
|
|
||||||
for name, u := range kept {
|
|
||||||
hashes[name] = u.PasswordHash
|
|
||||||
}
|
|
||||||
_, missing := broker.WithPasswords(users, hashes)
|
|
||||||
wanted := map[string]bool{}
|
|
||||||
for _, m := range missing {
|
|
||||||
wanted[m] = true
|
|
||||||
}
|
|
||||||
|
|
||||||
var machines, modules, skipped int
|
|
||||||
for _, p := range users {
|
|
||||||
if !again && !wanted[p.Username()] {
|
|
||||||
continue
|
|
||||||
}
|
|
||||||
switch p.Kind {
|
|
||||||
case broker.KindController:
|
|
||||||
if controllerNode == "" {
|
|
||||||
return errors.New("the control plane is not assigned anywhere, so its credential has nowhere to go")
|
|
||||||
}
|
|
||||||
password, err := inv.MintBusPassword(ctx, inventory.BusUser{Username: p.Username(), Kind: inventory.BusController})
|
|
||||||
if err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
url := "nats://" + p.Username() + ":" + password + "@" + busAddress
|
|
||||||
if err := inv.AcceptSecretForModule(ctx, controllerNode, "mesh-controller", "bus", url); err != nil {
|
|
||||||
return fmt.Errorf("the control plane's credential is minted and could not be sealed to %s: %w", controllerNode, err)
|
|
||||||
}
|
|
||||||
fmt.Printf("control plane: credential minted, sealed to %s as its `bus` secret\n", controllerNode)
|
|
||||||
|
|
||||||
case broker.KindNode:
|
|
||||||
password, err := inv.MintBusPassword(ctx, inventory.BusUser{Username: p.Username(), Kind: inventory.BusNode, Node: p.Node})
|
|
||||||
if err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
membership, _ := json.Marshal(map[string]string{
|
|
||||||
"broker": busAddress, "fingerprint": known.Fingerprint, "password": password, "transport": "nats",
|
|
||||||
})
|
|
||||||
key, err := inv.SealingKeyOf(ctx, p.Node)
|
|
||||||
if err != nil {
|
|
||||||
return fmt.Errorf("%s has no sealing key, so its membership cannot be sealed to it: %w", p.Node, err)
|
|
||||||
}
|
|
||||||
sealed, err := secrets.Seal(key, membership)
|
|
||||||
if err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
if err := inv.PutBusMembership(ctx, p.Node, sealed); err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
machines++
|
|
||||||
|
|
||||||
case broker.KindModule:
|
|
||||||
if p.Module == "mesh-controller" {
|
|
||||||
// The control plane is a module too, and its `broker` secret is the old bus's
|
|
||||||
// credential it is still using while this runs. Writing the new bus's blob there
|
|
||||||
// cut the mesh off from its own old bus mid-move (2026-09-28). Its new-bus credential
|
|
||||||
// is the controller principal's `bus` secret above; nothing else is needed here.
|
|
||||||
skipped++
|
|
||||||
continue
|
|
||||||
}
|
|
||||||
m, inShelf := shelf[p.Module]
|
|
||||||
if !inShelf {
|
|
||||||
skipped++
|
|
||||||
continue
|
|
||||||
}
|
|
||||||
if _, reads := m.OwnSecrets["broker"]; !reads {
|
|
||||||
fmt.Printf(" %s on %s speaks on the bus but declares no `broker` secret to receive a credential in; skipped\n", p.Module, p.Node)
|
|
||||||
skipped++
|
|
||||||
continue
|
|
||||||
}
|
|
||||||
password, err := inv.MintBusPassword(ctx, inventory.BusUser{Username: p.Username(), Kind: inventory.BusModule, Node: p.Node, Module: p.Module})
|
|
||||||
if err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
if err := issueWith(ctx, inv, m, p.Node, "", known, busAddress, p.Username(), password); err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
modules++
|
|
||||||
|
|
||||||
default:
|
|
||||||
skipped++
|
|
||||||
}
|
|
||||||
}
|
|
||||||
fmt.Printf("minted for %d machine(s) and %d module runtime(s); %d skipped; the bus is at %s\n",
|
|
||||||
machines, modules, skipped, busAddress)
|
|
||||||
fmt.Println(" each machine's membership and each module's credential arrive with the next push of its machine;")
|
|
||||||
fmt.Println(" push the machine running the bus first, so the bus stands with its user list before anything dials it")
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
|
|
||||||
// providesBus is whether a manifest provides the mesh's bus.
|
|
||||||
func providesBus(m catalogue.Manifest) bool {
|
|
||||||
for _, o := range m.Provides {
|
|
||||||
if o.Name == "mesh-bus" {
|
|
||||||
return true
|
|
||||||
}
|
|
||||||
}
|
|
||||||
return false
|
|
||||||
}
|
|
||||||
|
|||||||
@@ -1,14 +1,8 @@
|
|||||||
package broker
|
package broker
|
||||||
|
|
||||||
import (
|
import (
|
||||||
"crypto/sha256"
|
|
||||||
"crypto/tls"
|
|
||||||
"crypto/x509"
|
|
||||||
"encoding/hex"
|
|
||||||
"errors"
|
"errors"
|
||||||
"fmt"
|
"fmt"
|
||||||
"os"
|
|
||||||
"strings"
|
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/nats-io/nats.go"
|
"github.com/nats-io/nats.go"
|
||||||
@@ -30,58 +24,21 @@ type JetStream struct {
|
|||||||
|
|
||||||
// Dial connects and returns the controller's JetStream handle.
|
// Dial connects and returns the controller's JetStream handle.
|
||||||
func Dial(url string, opts ...nats.Option) (*JetStream, error) {
|
func Dial(url string, opts ...nats.Option) (*JetStream, error) {
|
||||||
|
// A name, because a connection nobody can identify in the server's own monitoring is one
|
||||||
|
// nobody can attribute a problem to.
|
||||||
opts = append(opts, nats.Name("mesh-controller"), nats.Timeout(10*time.Second))
|
opts = append(opts, nats.Name("mesh-controller"), nats.Timeout(10*time.Second))
|
||||||
// **Pinned, not named.** The bus presents the mesh's own certificate, which names nothing a
|
|
||||||
// public verifier would accept (design 25 §4: a host pins the server's exact certificate and
|
|
||||||
// checks nothing else, and so does this). Without this, the first connection failed with
|
|
||||||
// "certificate is not valid for any names" against a bus that was answering (2026-09-28).
|
|
||||||
if path := strings.TrimSpace(os.Getenv(CertificateVar)); path != "" {
|
|
||||||
pinned, err := pinnedTo(path)
|
|
||||||
if err != nil {
|
|
||||||
return nil, err
|
|
||||||
}
|
|
||||||
opts = append(opts, nats.Secure(pinned))
|
|
||||||
}
|
|
||||||
// The address in an error is the address alone. The URL carries this controller's password,
|
|
||||||
// and an error here is written on the assumption it will be logged.
|
|
||||||
where := BareAddress(url)
|
|
||||||
conn, err := nats.Connect(url, opts...)
|
conn, err := nats.Connect(url, opts...)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, fmt.Errorf("connecting to the bus at %s: %w", where, err)
|
return nil, fmt.Errorf("connecting to the bus at %s: %w", url, err)
|
||||||
}
|
}
|
||||||
js, err := conn.JetStream()
|
js, err := conn.JetStream()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
conn.Close()
|
conn.Close()
|
||||||
return nil, fmt.Errorf("the bus at %s has no JetStream: %w", where, err)
|
return nil, fmt.Errorf("the bus at %s has no JetStream: %w", url, err)
|
||||||
}
|
}
|
||||||
return &JetStream{conn: conn, js: js}, nil
|
return &JetStream{conn: conn, js: js}, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// pinnedTo is a TLS configuration that accepts exactly the certificate in the file and no other:
|
|
||||||
// the leaf's SHA-256, compared on every handshake, with the name and the chain deliberately not
|
|
||||||
// consulted — a self-signed certificate with no names is the ordinary case for a mesh's bus.
|
|
||||||
func pinnedTo(path string) (*tls.Config, error) {
|
|
||||||
want, err := FingerprintOf(path)
|
|
||||||
if err != nil {
|
|
||||||
return nil, err
|
|
||||||
}
|
|
||||||
return &tls.Config{
|
|
||||||
InsecureSkipVerify: true, //nolint:gosec // replaced by the pin below, which is stricter
|
|
||||||
MinVersion: tls.VersionTLS12,
|
|
||||||
VerifyPeerCertificate: func(rawCerts [][]byte, _ [][]*x509.Certificate) error {
|
|
||||||
if len(rawCerts) == 0 {
|
|
||||||
return errors.New("the bus presented no certificate")
|
|
||||||
}
|
|
||||||
sum := sha256.Sum256(rawCerts[0])
|
|
||||||
got := "sha256:" + hex.EncodeToString(sum[:])
|
|
||||||
if got != want {
|
|
||||||
return fmt.Errorf("the bus presented a certificate this mesh does not know (%s…), expected %s…", got[:23], want[:23])
|
|
||||||
}
|
|
||||||
return nil
|
|
||||||
},
|
|
||||||
}, nil
|
|
||||||
}
|
|
||||||
|
|
||||||
// Conn is the connection itself, for what the mesh keeps off JetStream on purpose — a heartbeat,
|
// Conn is the connection itself, for what the mesh keeps off JetStream on purpose — a heartbeat,
|
||||||
// a tool call — where a lost message is answered by the next one or by a timeout the caller
|
// a tool call — where a lost message is answered by the next one or by a timeout the caller
|
||||||
// already handles (design 25 §3).
|
// already handles (design 25 §3).
|
||||||
|
|||||||
@@ -244,14 +244,7 @@ func PermissionsFor(p Principal) (Permissions, error) {
|
|||||||
case KindNode:
|
case KindNode:
|
||||||
// A host publishes its own node's control traffic and subscribes its own declaration —
|
// A host publishes its own node's control traffic and subscribes its own declaration —
|
||||||
// and nothing of any other node's.
|
// and nothing of any other node's.
|
||||||
// And binding to its consumer, which asks the server about it (CONSUMER.INFO) — the one
|
pub = []string{"mesh.control." + p.Node + ".>"}
|
||||||
// thing the host does that nothing granted. Found the first time a machine dialled a
|
|
||||||
// permissioned server: "this node cannot read its declarations" (2026-09-28). The ack and
|
|
||||||
// the inbox are granted below with every principal's.
|
|
||||||
pub = []string{
|
|
||||||
"mesh.control." + p.Node + ".>",
|
|
||||||
"$JS.API.CONSUMER.INFO.NODES." + p.Node,
|
|
||||||
}
|
|
||||||
sub = []string{"mesh.node." + p.Node + ".declare"}
|
sub = []string{"mesh.node." + p.Node + ".declare"}
|
||||||
|
|
||||||
case KindModule:
|
case KindModule:
|
||||||
|
|||||||
+1
-1
@@ -32,7 +32,7 @@ accounts {
|
|||||||
subscribe: { allow: ["_INBOX.enrol.one.>"] }
|
subscribe: { allow: ["_INBOX.enrol.one.>"] }
|
||||||
} }
|
} }
|
||||||
{ user: "node.one", password: "$2a$11$nnnnnnnnnnnnnnnnnnnnnn", permissions: {
|
{ user: "node.one", password: "$2a$11$nnnnnnnnnnnnnnnnnnnnnn", permissions: {
|
||||||
publish: { allow: ["$JS.ACK.NODES.one.>", "$JS.API.CONSUMER.INFO.NODES.one", "mesh.control.one.>"] }
|
publish: { allow: ["$JS.ACK.NODES.one.>", "mesh.control.one.>"] }
|
||||||
subscribe: { allow: ["_INBOX.node.one.>", "mesh.node.one.declare"] }
|
subscribe: { allow: ["_INBOX.node.one.>", "mesh.node.one.declare"] }
|
||||||
} }
|
} }
|
||||||
{ user: "one.telegram", password: "$2a$11$tttttttttttttttttttttt", permissions: {
|
{ user: "one.telegram", password: "$2a$11$tttttttttttttttttttttt", permissions: {
|
||||||
|
|||||||
@@ -103,11 +103,6 @@ type Rendering struct {
|
|||||||
// **Only the users, never the server's own settings**: those are the module's, in its image and
|
// **Only the users, never the server's own settings**: those are the module's, in its image and
|
||||||
// its mounts (Manifest.BusUsers).
|
// its mounts (Manifest.BusUsers).
|
||||||
BusUsers string
|
BusUsers string
|
||||||
// BusMembership is this machine's membership for the bus the mesh is moving to, sealed to it
|
|
||||||
// (design 28, task 5.2). Empty for a machine not being moved. Written as a file the host reads
|
|
||||||
// after the declaration has applied, so the bus it names is standing before the machine leaves
|
|
||||||
// the one it is on.
|
|
||||||
BusMembership string
|
|
||||||
|
|
||||||
// MeshRange is the private network's CIDR (the range node addresses are allocated from), for a
|
// MeshRange is the private network's CIDR (the range node addresses are allocated from), for a
|
||||||
// module that must name the whole mesh rather than one machine — an intrusion filter that must
|
// module that must name the whole mesh rather than one machine — an intrusion filter that must
|
||||||
@@ -243,23 +238,9 @@ func (r Resolution) Compose(with Rendering) (Composed, error) {
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
return Composed{}, err
|
return Composed{}, err
|
||||||
}
|
}
|
||||||
if with.BusMembership != "" {
|
|
||||||
// The machine's own, not any module's: how it reaches the mesh from now on. Sealed like a
|
|
||||||
// secret and placed where the host looks for exactly this (design 28, task 5.2).
|
|
||||||
resources = append(resources, map[string]any{
|
|
||||||
"id": BusMembershipID(), "type": "file", "path": BusMembershipPath,
|
|
||||||
"sealed": with.BusMembership, "mode": "0600",
|
|
||||||
})
|
|
||||||
}
|
|
||||||
return Composed{Resources: resources, Owner: owner}, nil
|
return Composed{Resources: resources, Owner: owner}, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// BusMembershipID names the resource carrying a machine's membership for the new bus, and
|
|
||||||
// BusMembershipPath is where the host reads it — the same constant on both sides.
|
|
||||||
func BusMembershipID() string { return "bus-membership" }
|
|
||||||
|
|
||||||
const BusMembershipPath = "/var/lib/mesh/membership-next.json"
|
|
||||||
|
|
||||||
func (r Resolution) compose(with Rendering, owner map[string]string) ([]map[string]any, error) {
|
func (r Resolution) compose(with Rendering, owner map[string]string) ([]map[string]any, error) {
|
||||||
// Every manifest is placed first (novox/hq ADR 0112): the maps naming where its bindings,
|
// Every manifest is placed first (novox/hq ADR 0112): the maps naming where its bindings,
|
||||||
// credentials and contributions land are resolved against this node's directories, so every
|
// credentials and contributions land are resolved against this node's directories, so every
|
||||||
|
|||||||
@@ -114,32 +114,3 @@ func TestCanHoldJudgesClaimScopeAndWhatTheSeatDelivers(t *testing.T) {
|
|||||||
t.Fatalf("with the row saying amqp, an amqp provider was refused: %v", err)
|
t.Fatalf("with the row saying amqp, an amqp provider was refused: %v", err)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// A machine being moved is handed its membership for the new bus as a sealed file in its own
|
|
||||||
// declaration — the machine's, not any module's (design 28, task 5.2).
|
|
||||||
func TestAMembershipForTheNewBusIsComposedAsASealedFile(t *testing.T) {
|
|
||||||
r := Resolution{Node: "anchor"}
|
|
||||||
got, err := r.Compose(Rendering{BusMembership: "sealed-blob"})
|
|
||||||
if err != nil {
|
|
||||||
t.Fatal(err)
|
|
||||||
}
|
|
||||||
var found map[string]any
|
|
||||||
for _, res := range got.Resources {
|
|
||||||
if res["id"] == BusMembershipID() {
|
|
||||||
found = res
|
|
||||||
}
|
|
||||||
}
|
|
||||||
if found == nil {
|
|
||||||
t.Fatalf("no membership resource in %v", got.Resources)
|
|
||||||
}
|
|
||||||
if found["path"] != BusMembershipPath || found["sealed"] != "sealed-blob" || found["mode"] != "0600" {
|
|
||||||
t.Fatalf("the membership is not a sealed 0600 file where the host reads it: %v", found)
|
|
||||||
}
|
|
||||||
// And a machine not being moved is handed nothing.
|
|
||||||
got, _ = r.Compose(Rendering{})
|
|
||||||
for _, res := range got.Resources {
|
|
||||||
if res["id"] == BusMembershipID() {
|
|
||||||
t.Fatal("a machine with no membership on record was handed one")
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|||||||
@@ -230,35 +230,3 @@ func (i *Inventory) ForgetPerson(ctx context.Context, name string) error {
|
|||||||
}
|
}
|
||||||
return i.ForgetBusUser(ctx, "person."+name)
|
return i.ForgetBusUser(ctx, "person."+name)
|
||||||
}
|
}
|
||||||
|
|
||||||
// PutBusMembership records a machine's membership for the new bus, sealed to it (design 28, 5.2).
|
|
||||||
// Replaces any earlier one: a machine has one membership per bus, and re-minting is re-telling.
|
|
||||||
func (i *Inventory) PutBusMembership(ctx context.Context, nodeName, sealed string) error {
|
|
||||||
node, err := i.NodeByName(ctx, nodeName)
|
|
||||||
if err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
_, err = i.store.Pool().Exec(ctx,
|
|
||||||
`insert into bus_membership (node, sealed) values ($1, $2)
|
|
||||||
on conflict (node) do update set sealed = excluded.sealed, since = now()`, node.ID, sealed)
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
|
|
||||||
// BusMemberships is every machine's sealed membership for the new bus, by node name.
|
|
||||||
func (i *Inventory) BusMemberships(ctx context.Context) (map[string]string, error) {
|
|
||||||
rows, err := i.store.Pool().Query(ctx,
|
|
||||||
`select n.name, b.sealed from bus_membership b join node n on n.id = b.node`)
|
|
||||||
if err != nil {
|
|
||||||
return nil, err
|
|
||||||
}
|
|
||||||
defer rows.Close()
|
|
||||||
out := map[string]string{}
|
|
||||||
for rows.Next() {
|
|
||||||
var name, sealed string
|
|
||||||
if err := rows.Scan(&name, &sealed); err != nil {
|
|
||||||
return nil, err
|
|
||||||
}
|
|
||||||
out[name] = sealed
|
|
||||||
}
|
|
||||||
return out, rows.Err()
|
|
||||||
}
|
|
||||||
|
|||||||
@@ -80,20 +80,3 @@ func TestUnassigningTheHolderTakesTheHoldingWithIt(t *testing.T) {
|
|||||||
t.Fatalf("the holding outlived the assignment it pointed at: %+v", held)
|
t.Fatalf("the holding outlived the assignment it pointed at: %+v", held)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
func TestAMachinesMembershipIsOneRowReplacedAndGoesWithTheMachine(t *testing.T) {
|
|
||||||
inv, ctx := twoBrokersOnTwoNodes(t)
|
|
||||||
if err := inv.PutBusMembership(ctx, "anchor", "first"); err != nil {
|
|
||||||
t.Fatal(err)
|
|
||||||
}
|
|
||||||
if err := inv.PutBusMembership(ctx, "anchor", "second"); err != nil {
|
|
||||||
t.Fatal(err)
|
|
||||||
}
|
|
||||||
got, err := inv.BusMemberships(ctx)
|
|
||||||
if err != nil {
|
|
||||||
t.Fatal(err)
|
|
||||||
}
|
|
||||||
if got["anchor"] != "second" || len(got) != 1 {
|
|
||||||
t.Fatalf("a re-told membership did not replace the first: %v", got)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|||||||
-11
@@ -1,11 +0,0 @@
|
|||||||
-- A machine already enrolled is moved to the new bus by being told its membership for it
|
|
||||||
-- (novox/hq design 28, task 5.2). Until this, a membership — bus address, fingerprint, password,
|
|
||||||
-- transport — existed only in the enrolment reply, and nothing could hand one to a machine that
|
|
||||||
-- had already joined. The row is the membership sealed to that machine, composed into its
|
|
||||||
-- declaration as a file it reads after applying; the plaintext exists once, at minting, and then
|
|
||||||
-- only on the machine. One per node: the mesh moves to one bus.
|
|
||||||
create table bus_membership (
|
|
||||||
node uuid primary key references node(id) on delete cascade,
|
|
||||||
sealed text not null,
|
|
||||||
since timestamptz not null default now()
|
|
||||||
);
|
|
||||||
@@ -1,73 +0,0 @@
|
|||||||
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},
|
|
||||||
js: js,
|
|
||||||
enroller: enroller,
|
|
||||||
listener: listener,
|
|
||||||
log: newLog(),
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// Bus is the controller's outbound, whichever transport it connected over. Callers that send a
|
|
||||||
// declaration or ask a tool use this rather than the channel, which one transport does not have.
|
|
||||||
func (s *Server) Bus() Bus { return s.bus }
|
|
||||||
@@ -7,7 +7,6 @@ import (
|
|||||||
"encoding/json"
|
"encoding/json"
|
||||||
"errors"
|
"errors"
|
||||||
"fmt"
|
"fmt"
|
||||||
"github.com/novox/mesh-controller/internal/broker"
|
|
||||||
"log"
|
"log"
|
||||||
"os"
|
"os"
|
||||||
"time"
|
"time"
|
||||||
@@ -71,7 +70,6 @@ type Server struct {
|
|||||||
bus Bus
|
bus Bus
|
||||||
conn *amqp.Connection
|
conn *amqp.Connection
|
||||||
channel *amqp.Channel
|
channel *amqp.Channel
|
||||||
js *broker.JetStream
|
|
||||||
|
|
||||||
enroller Enroller
|
enroller Enroller
|
||||||
listener Listener
|
listener Listener
|
||||||
@@ -182,12 +180,10 @@ func Connect(enroller Enroller, listener Listener) (*Server, error) {
|
|||||||
channel: channel,
|
channel: channel,
|
||||||
enroller: enroller,
|
enroller: enroller,
|
||||||
listener: listener,
|
listener: listener,
|
||||||
log: newLog(),
|
log: log.New(os.Stdout, "", log.LstdFlags),
|
||||||
}, nil
|
}, 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.
|
// 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) Channel() *amqp.Channel { return s.channel }
|
||||||
|
|
||||||
@@ -201,9 +197,6 @@ func (s *Server) Close() {
|
|||||||
if s.conn != nil {
|
if s.conn != nil {
|
||||||
_ = s.conn.Close()
|
_ = s.conn.Close()
|
||||||
}
|
}
|
||||||
if s.js != nil {
|
|
||||||
s.js.Close()
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// Serve acts on what arrives until the context ends.
|
// Serve acts on what arrives until the context ends.
|
||||||
|
|||||||
Executable
BIN
Binary file not shown.
+4
-5
@@ -23,8 +23,7 @@
|
|||||||
"licences": "/var/lib/mesh/mesh-controller/licences",
|
"licences": "/var/lib/mesh/mesh-controller/licences",
|
||||||
"broker": "/var/lib/mesh/mesh-controller/broker",
|
"broker": "/var/lib/mesh/mesh-controller/broker",
|
||||||
"broker-management": "/var/lib/mesh/mesh-controller/broker-management",
|
"broker-management": "/var/lib/mesh/mesh-controller/broker-management",
|
||||||
"broker-address": "/var/lib/mesh/mesh-controller/broker-address",
|
"broker-address": "/var/lib/mesh/mesh-controller/broker-address"
|
||||||
"bus": "/var/lib/mesh/mesh-controller/bus"
|
|
||||||
},
|
},
|
||||||
"secrets-owner": "65534:65534",
|
"secrets-owner": "65534:65534",
|
||||||
"resources": [
|
"resources": [
|
||||||
@@ -47,14 +46,15 @@
|
|||||||
"MESH_STORE_INVENTORY_FILE": "/run/secrets/inventory",
|
"MESH_STORE_INVENTORY_FILE": "/run/secrets/inventory",
|
||||||
"MESH_STORE_IDENTITY_FILE": "/run/secrets/identity",
|
"MESH_STORE_IDENTITY_FILE": "/run/secrets/identity",
|
||||||
"MESH_STORE_LICENCES_FILE": "/run/secrets/licences",
|
"MESH_STORE_LICENCES_FILE": "/run/secrets/licences",
|
||||||
|
"MESH_BROKER_AMQP_FILE": "/run/secrets/broker",
|
||||||
"MESH_BROKER_MANAGEMENT_FILE": "/run/secrets/broker-management",
|
"MESH_BROKER_MANAGEMENT_FILE": "/run/secrets/broker-management",
|
||||||
"MESH_BROKER_ADDRESS_FILE": "/run/secrets/broker-address",
|
"MESH_BROKER_ADDRESS_FILE": "/run/secrets/broker-address",
|
||||||
"MESH_STORE_INVENTORY_PORT": "${seat:mesh-store:5432}",
|
"MESH_STORE_INVENTORY_PORT": "${seat:mesh-store:5432}",
|
||||||
"MESH_STORE_IDENTITY_PORT": "${seat:mesh-store:5432}",
|
"MESH_STORE_IDENTITY_PORT": "${seat:mesh-store:5432}",
|
||||||
"MESH_STORE_LICENCES_PORT": "${seat:mesh-store:5432}",
|
"MESH_STORE_LICENCES_PORT": "${seat:mesh-store:5432}",
|
||||||
|
"MESH_BROKER_AMQP_PORT": "${seat:mesh-broker:5672}",
|
||||||
"MESH_BROKER_MANAGEMENT_PORT": "${seat:mesh-broker:15672}",
|
"MESH_BROKER_MANAGEMENT_PORT": "${seat:mesh-broker:15672}",
|
||||||
"MESH_BROKER_ADDRESS_PORT": "${seat:mesh-broker:5671}",
|
"MESH_BROKER_ADDRESS_PORT": "${seat:mesh-broker:5671}"
|
||||||
"MESH_BUS_NATS_FILE": "/run/secrets/bus"
|
|
||||||
},
|
},
|
||||||
"volumes": [
|
"volumes": [
|
||||||
"/var/lib/mesh-broker-tls:/broker-tls:ro",
|
"/var/lib/mesh-broker-tls:/broker-tls:ro",
|
||||||
@@ -62,7 +62,6 @@
|
|||||||
"/var/lib/mesh/mesh-controller/identity:/run/secrets/identity:ro",
|
"/var/lib/mesh/mesh-controller/identity:/run/secrets/identity:ro",
|
||||||
"/var/lib/mesh/mesh-controller/licences:/run/secrets/licences:ro",
|
"/var/lib/mesh/mesh-controller/licences:/run/secrets/licences:ro",
|
||||||
"/var/lib/mesh/mesh-controller/broker:/run/secrets/broker:ro",
|
"/var/lib/mesh/mesh-controller/broker:/run/secrets/broker:ro",
|
||||||
"/var/lib/mesh/mesh-controller/bus:/run/secrets/bus:ro",
|
|
||||||
"/var/lib/mesh/mesh-controller/broker-management:/run/secrets/broker-management:ro",
|
"/var/lib/mesh/mesh-controller/broker-management:/run/secrets/broker-management:ro",
|
||||||
"/var/lib/mesh/mesh-controller/broker-address:/run/secrets/broker-address:ro"
|
"/var/lib/mesh/mesh-controller/broker-address:/run/secrets/broker-address:ro"
|
||||||
],
|
],
|
||||||
|
|||||||
Reference in New Issue
Block a user