Unify trunk on main: initialization → main #3
+36
-5
@@ -222,7 +222,7 @@ func run(ctx context.Context, command string, opts options) error {
|
||||
return w.Flush()
|
||||
|
||||
case "enrol", "enroll":
|
||||
return enrol(opts)
|
||||
return enrol(ctx, opts)
|
||||
|
||||
case "version":
|
||||
fmt.Println(version)
|
||||
@@ -386,7 +386,7 @@ func runApply(ctx context.Context, opts options, d *declaration.Declaration, sou
|
||||
// the one-time secret together with a public key it generated itself.
|
||||
//
|
||||
// The mesh issues no identity. This machine arrives holding one; what it receives is being known.
|
||||
func enrol(opts options) error {
|
||||
func enrol(ctx context.Context, opts options) error {
|
||||
tokenText, name := &opts.token, &opts.nodeName
|
||||
if strings.TrimSpace(*tokenText) == "" {
|
||||
return errors.New("enrol --token <token>: the token is carried to this machine by a " +
|
||||
@@ -429,13 +429,44 @@ func enrol(opts options) error {
|
||||
defer conn.Close()
|
||||
fmt.Println("\nthe broker presented the certificate this token pins")
|
||||
|
||||
conn.Close()
|
||||
|
||||
mine, err := identity.Generate(*name)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
fmt.Printf("generated this node's identity: %s\n", mine.PublicBase64())
|
||||
|
||||
return errors.New("the link is not built: this machine has verified the broker and made its " +
|
||||
"identity, and there is nothing yet to present them to.\n" +
|
||||
"Nothing has been saved, so this can be run again unchanged")
|
||||
// What this machine can be asked to do, gathered before joining rather than after. The
|
||||
// control plane cannot decide what a node should run without it, so it travels with the
|
||||
// request instead of being asked for in a second round trip.
|
||||
detected := profile.Detect(ctx, profile.Default(nil), opts.timeout)
|
||||
reported := map[string]any{}
|
||||
if raw, err := json.Marshal(detected); err == nil {
|
||||
_ = json.Unmarshal(raw, &reported)
|
||||
}
|
||||
|
||||
reply, err := link.Enrol(ctx, token.Broker, token.Fingerprint, *name, token.Secret,
|
||||
mine.Public, reported, opts.timeout)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
// The mesh's name for this node wins over what the machine called itself: the token was
|
||||
// issued for a node record, and that record is what the identity binds to.
|
||||
mine.Node = reply.Node
|
||||
|
||||
// Saved only now, and only once the mesh has said it knows this node. A node holding an
|
||||
// identity the mesh has never recorded would believe it had joined and be believed by
|
||||
// nobody — worse than not having joined, because nothing would look wrong.
|
||||
if err := identity.Save(identityPath, mine); err != nil {
|
||||
return fmt.Errorf(
|
||||
"the mesh accepted this node as %q and its identity could not be saved: %w\n"+
|
||||
"That token is spent, so getting back needs a new one", reply.Node, err)
|
||||
}
|
||||
|
||||
fmt.Printf("\nenrolled as %s\n", reply.Node)
|
||||
fmt.Printf(" identity %s\n", identityPath)
|
||||
fmt.Printf(" queue %s\n", reply.Queue)
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -1,11 +1,14 @@
|
||||
// substrate-first-node.lock — what a machine must be before a mesh exists.
|
||||
//
|
||||
// Steps 0 to 4 of the bootstrap (novox/hq 03-DESIGN/01-to-be/07-the-substrate.md): a container
|
||||
// Steps 0 to 5 of the bootstrap (novox/hq 03-DESIGN/01-to-be/07-the-substrate.md): a container
|
||||
// runtime, a store, a database per context, that context's schema, and the broker.
|
||||
//
|
||||
// It stops before step 5 (a virtual host, a credential, a certificate) and step 6 (the control
|
||||
// plane runs), because nothing consumes them yet. A bundle naming a control plane that serves
|
||||
// nothing would be a bundle whose last step cannot be checked.
|
||||
// It stops before step 6, where the control plane runs.
|
||||
//
|
||||
// The broker generates its OWN certificate, in its own image, into a volume it then mounts read
|
||||
// only. Self-signed, because at this moment there is no mesh to issue one and no public name to
|
||||
// obtain one for -- and it does not matter, because what a joining node checks is the fingerprint
|
||||
// pinned in its token, not a chain or a name (novox/hq ADR 0004). The subject is decoration.
|
||||
//
|
||||
// PINNED BY DIGEST, and the digest is not decoration: a tag can be made to point at a different
|
||||
// image, and this file is applied on a machine with no mesh to ask about anything. These digests
|
||||
@@ -41,6 +44,7 @@
|
||||
"POSTGRES_PASSWORD": "bootstrap",
|
||||
"PGDATA": "/var/lib/postgresql/data/pgdata"
|
||||
},
|
||||
"ports": ["127.0.0.1:5432:5432"],
|
||||
"volumes": ["mesh-store-data:/var/lib/postgresql/data"]
|
||||
},
|
||||
{
|
||||
@@ -58,21 +62,40 @@
|
||||
"verify": ["sh", "-c", "psql -U postgres -lqt | cut -d'|' -f1 | grep -qw inventory"]
|
||||
},
|
||||
{
|
||||
"id": "inventory-schema",
|
||||
"id": "identity-database",
|
||||
"type": "action",
|
||||
"in": "mesh-store",
|
||||
"command": ["sh", "-c", "psql -U postgres -c 'CREATE DATABASE identity'"],
|
||||
"verify": ["sh", "-c", "psql -U postgres -lqt | cut -d'|' -f1 | grep -qw identity"]
|
||||
},
|
||||
{
|
||||
"id": "context-schemas",
|
||||
"type": "action",
|
||||
"command": ["docker", "run", "--rm", "--network", "container:mesh-store",
|
||||
"-e", "MESH_STORE_INVENTORY=postgres://postgres:bootstrap@127.0.0.1:5432/inventory?sslmode=disable",
|
||||
"192.0.2.250:5000/mesh-control@sha256:1c27a43c4431c2b580404e8e1768cd858009e265e80b6f1591eb6de1123fc411",
|
||||
"-e", "MESH_STORE_IDENTITY=postgres://postgres:bootstrap@127.0.0.1:5432/identity?sslmode=disable",
|
||||
"192.0.2.250:5000/mesh-control@sha256:c6e96dc574ea52085bae4ed0a64e593ee512265ecb226643aa1fcbaf56d78396",
|
||||
"migrate"],
|
||||
"verify": ["sh", "-c", "docker exec mesh-store psql -U postgres -d inventory -tAc \"select to_regclass('public.node')\" | grep -qx node"]
|
||||
"verify": ["sh", "-c", "docker exec mesh-store psql -U postgres -d inventory -tAc \"select to_regclass('public.node')\" | grep -qx node && docker exec mesh-store psql -U postgres -d identity -tAc \"select to_regclass('public.signing_key')\" | grep -qx signing_key"]
|
||||
},
|
||||
{
|
||||
"id": "broker-certificate",
|
||||
"type": "action",
|
||||
"command": ["docker", "run", "--rm", "--entrypoint", "sh", "-v", "mesh-broker-tls:/tls",
|
||||
"192.0.2.250:5000/cloudamqp/lavinmq@sha256:b117c254e6e269a29db479e6b410ca4e46e035b4981e49d24b159673ef09d336",
|
||||
"-c", "test -f /tls/tls.crt || (openssl req -x509 -newkey rsa:2048 -nodes -keyout /tls/tls.key -out /tls/tls.crt -days 3650 -subj '/CN=mesh-broker' >/dev/null 2>&1 && chmod 644 /tls/tls.crt && chmod 600 /tls/tls.key)"],
|
||||
"verify": ["docker", "run", "--rm", "--entrypoint", "sh", "-v", "mesh-broker-tls:/tls",
|
||||
"192.0.2.250:5000/cloudamqp/lavinmq@sha256:b117c254e6e269a29db479e6b410ca4e46e035b4981e49d24b159673ef09d336",
|
||||
"-c", "test -s /tls/tls.crt && openssl x509 -in /tls/tls.crt -noout"]
|
||||
},
|
||||
{
|
||||
"id": "broker",
|
||||
"type": "container",
|
||||
"name": "mesh-broker",
|
||||
"image": "192.0.2.250:5000/cloudamqp/lavinmq@sha256:b117c254e6e269a29db479e6b410ca4e46e035b4981e49d24b159673ef09d336",
|
||||
"ports": ["5672:5672"],
|
||||
"volumes": ["mesh-broker-data:/var/lib/lavinmq"]
|
||||
"ports": ["5671:5671", "127.0.0.1:5672:5672", "127.0.0.1:15672:15672"],
|
||||
"volumes": ["mesh-broker-data:/var/lib/lavinmq", "mesh-broker-tls:/tls:ro"],
|
||||
"args": ["--amqps-port=5671", "--cert=/tls/tls.crt", "--key=/tls/tls.key"]
|
||||
},
|
||||
{
|
||||
"id": "broker-ready",
|
||||
@@ -80,6 +103,24 @@
|
||||
"in": "mesh-broker",
|
||||
"command": ["sh", "-c", "for i in $(seq 1 60); do lavinmqctl status >/dev/null 2>&1 && exit 0; sleep 1; done; exit 1"],
|
||||
"verify": ["lavinmqctl", "status"]
|
||||
},
|
||||
{
|
||||
"id": "control-plane",
|
||||
"type": "container",
|
||||
"name": "mesh-control",
|
||||
"image": "192.0.2.250:5000/mesh-control@sha256:c6e96dc574ea52085bae4ed0a64e593ee512265ecb226643aa1fcbaf56d78396",
|
||||
"network": "host",
|
||||
"args": ["serve"],
|
||||
"volumes": ["mesh-broker-tls:/broker-tls:ro"],
|
||||
"env": {
|
||||
"MESH_STORE_INVENTORY": "postgres://postgres:bootstrap@127.0.0.1:5432/inventory?sslmode=disable",
|
||||
"MESH_STORE_IDENTITY": "postgres://postgres:bootstrap@127.0.0.1:5432/identity?sslmode=disable",
|
||||
"MESH_BROKER_AMQP": "amqp://guest:guest@127.0.0.1:5672/",
|
||||
"MESH_BROKER_MANAGEMENT": "http://guest:guest@127.0.0.1:15672",
|
||||
"MESH_BROKER_ADDRESS": "192.0.2.10:5671",
|
||||
"MESH_BROKER_CERTIFICATE": "/broker-tls/tls.crt"
|
||||
}
|
||||
}
|
||||
|
||||
]
|
||||
}
|
||||
|
||||
@@ -1,3 +1,5 @@
|
||||
module github.com/novox/mesh-host
|
||||
|
||||
go 1.24
|
||||
|
||||
require github.com/rabbitmq/amqp091-go v1.14.0 // indirect
|
||||
|
||||
@@ -0,0 +1,2 @@
|
||||
github.com/rabbitmq/amqp091-go v1.14.0 h1:RSaT7aOKt/OrkVUyswPDW29lnRz9psuGmfZFBmLqLek=
|
||||
github.com/rabbitmq/amqp091-go v1.14.0/go.mod h1:Hy4jKW5kQART1u+JkDTF9YYOQUHXqMuhrgxOEeS7G4o=
|
||||
@@ -584,8 +584,12 @@ func applyContainer(ctx context.Context, r *declaration.Container, run Runner) (
|
||||
}
|
||||
}
|
||||
|
||||
args := []string{"run", "--detach", "--name", r.Name, "--restart", "unless-stopped",
|
||||
"--label", specLabel + "=" + want, "--label", idLabel + "=" + r.ID}
|
||||
args := []string{"run", "--detach", "--name", r.Name, "--restart", "unless-stopped"}
|
||||
if r.Network != "" {
|
||||
args = append(args, "--network", r.Network)
|
||||
}
|
||||
args = append(args,
|
||||
"--label", specLabel+"="+want, "--label", idLabel+"="+r.ID)
|
||||
for _, k := range sortedKeys(r.Env) {
|
||||
args = append(args, "--env", k+"="+r.Env[k])
|
||||
}
|
||||
|
||||
@@ -168,6 +168,13 @@ type Container struct {
|
||||
Ports []string `json:"ports,omitempty"`
|
||||
Volumes []string `json:"volumes,omitempty"`
|
||||
Args []string `json:"args,omitempty"`
|
||||
// Network is the container's network, passed to the runtime unchanged.
|
||||
//
|
||||
// Needed because the control plane must reach the store and the broker on the machine it was
|
||||
// raised on, before there is any mesh to arrange that. The alternative was publishing ports
|
||||
// and guessing an address that works from inside a container, which is the same thing with a
|
||||
// worse failure mode.
|
||||
Network string `json:"network,omitempty"`
|
||||
}
|
||||
|
||||
func (c *Container) Identity() string { return c.ID }
|
||||
|
||||
@@ -0,0 +1,149 @@
|
||||
package link
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"net/url"
|
||||
"time"
|
||||
|
||||
amqp "github.com/rabbitmq/amqp091-go"
|
||||
)
|
||||
|
||||
// The wire format shared with the control plane, which defines it separately because this binary
|
||||
// requires nothing present and does not import it. A test on each side asserts the field names.
|
||||
const (
|
||||
Exchange = "mesh"
|
||||
KeyEnrol = "enrol"
|
||||
)
|
||||
|
||||
// QueueFor is the queue this node consumes from — the only one its account may read.
|
||||
func QueueFor(node string) string { return "node." + node }
|
||||
|
||||
// EnrolRequest is what this node says when joining.
|
||||
type EnrolRequest struct {
|
||||
Node string `json:"node"`
|
||||
Secret string `json:"secret"`
|
||||
PublicKey []byte `json:"public_key"`
|
||||
Profile map[string]any `json:"profile,omitempty"`
|
||||
}
|
||||
|
||||
// EnrolReply is what the mesh says back.
|
||||
type EnrolReply struct {
|
||||
Accepted bool `json:"accepted"`
|
||||
Node string `json:"node,omitempty"`
|
||||
Queue string `json:"queue,omitempty"`
|
||||
Refusal string `json:"refusal,omitempty"`
|
||||
}
|
||||
|
||||
// ErrRefused is what a node gets when the mesh will not have it.
|
||||
var ErrRefused = errors.New("the mesh refused this enrolment")
|
||||
|
||||
// Enrol presents this node's key and its one-time secret, and waits to be told it is known.
|
||||
//
|
||||
// The broker has already authenticated this connection: the account was created when the token
|
||||
// was issued and the secret is its password. So this is not how the node gets in — it is what it
|
||||
// says once it is in, and the secret travels again because the control plane must not have to ask
|
||||
// the broker who connected.
|
||||
func Enrol(ctx context.Context, address, pin, node, secret string, public []byte,
|
||||
profile map[string]any, timeout time.Duration) (EnrolReply, error) {
|
||||
|
||||
config, err := PinnedConfig(pin)
|
||||
if err != nil {
|
||||
return EnrolReply{}, err
|
||||
}
|
||||
|
||||
// The account name is the node's, and the password is the token's secret. Escaped because a
|
||||
// name or secret containing a colon or an at-sign would otherwise change which host this
|
||||
// connects to — a credential silently redirecting a connection is the worst shape this could
|
||||
// take.
|
||||
dsn := fmt.Sprintf("amqps://%s:%s@%s/",
|
||||
url.QueryEscape(node), url.QueryEscape(secret), address)
|
||||
|
||||
conn, err := amqp.DialConfig(dsn, amqp.Config{
|
||||
TLSClientConfig: config,
|
||||
Dial: amqp.DefaultDial(timeout),
|
||||
})
|
||||
if err != nil {
|
||||
if errors.Is(err, ErrWrongCertificate) {
|
||||
return EnrolReply{}, err
|
||||
}
|
||||
// Not quoted back: the DSN carries the one-time secret.
|
||||
return EnrolReply{}, fmt.Errorf("cannot reach the broker at %s as %s: %w", address, node, err)
|
||||
}
|
||||
defer conn.Close()
|
||||
|
||||
channel, err := conn.Channel()
|
||||
if err != nil {
|
||||
return EnrolReply{}, err
|
||||
}
|
||||
defer channel.Close()
|
||||
|
||||
// This node's own queue, which its account is scoped to and nothing else may read.
|
||||
queue, err := channel.QueueDeclare(QueueFor(node), true, false, false, false, nil)
|
||||
if err != nil {
|
||||
return EnrolReply{}, fmt.Errorf(
|
||||
"cannot declare this node's queue %s: %w", QueueFor(node), err)
|
||||
}
|
||||
|
||||
replies, err := channel.Consume(queue.Name, "", true, false, false, false, nil)
|
||||
if err != nil {
|
||||
return EnrolReply{}, err
|
||||
}
|
||||
|
||||
request := EnrolRequest{Node: node, Secret: secret, PublicKey: public, Profile: profile}
|
||||
body, err := json.Marshal(request)
|
||||
if err != nil {
|
||||
return EnrolReply{}, err
|
||||
}
|
||||
|
||||
correlation := fmt.Sprintf("%s-%d", node, time.Now().UnixNano())
|
||||
publish, cancel := context.WithTimeout(ctx, timeout)
|
||||
defer cancel()
|
||||
if err := channel.PublishWithContext(publish, Exchange, KeyEnrol, false, false,
|
||||
amqp.Publishing{
|
||||
ContentType: "application/json",
|
||||
CorrelationId: correlation,
|
||||
ReplyTo: queue.Name,
|
||||
Body: body,
|
||||
}); err != nil {
|
||||
return EnrolReply{}, fmt.Errorf("cannot publish to the %s exchange: %w", Exchange, err)
|
||||
}
|
||||
|
||||
// Waited for rather than assumed. A published message that nothing answers means the control
|
||||
// plane is not running, and a node that carried on regardless would believe it had joined a
|
||||
// mesh that has never heard of it.
|
||||
deadline := time.NewTimer(timeout)
|
||||
defer deadline.Stop()
|
||||
closed := conn.NotifyClose(make(chan *amqp.Error, 1))
|
||||
|
||||
for {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return EnrolReply{}, ctx.Err()
|
||||
case reason := <-closed:
|
||||
return EnrolReply{}, fmt.Errorf("the broker closed the connection: %v", reason)
|
||||
case <-deadline.C:
|
||||
return EnrolReply{}, fmt.Errorf(
|
||||
"the broker accepted this node's connection and nothing answered within %s. The "+
|
||||
"mesh's broker is running and its control plane is not", timeout)
|
||||
case delivery, ok := <-replies:
|
||||
if !ok {
|
||||
return EnrolReply{}, errors.New("the broker stopped delivering")
|
||||
}
|
||||
// Anything else on this queue is not the answer to this question.
|
||||
if delivery.CorrelationId != correlation {
|
||||
continue
|
||||
}
|
||||
var reply EnrolReply
|
||||
if err := json.Unmarshal(delivery.Body, &reply); err != nil {
|
||||
return EnrolReply{}, fmt.Errorf("the mesh's answer could not be read: %w", err)
|
||||
}
|
||||
if !reply.Accepted {
|
||||
return reply, fmt.Errorf("%w: %s", ErrRefused, reply.Refusal)
|
||||
}
|
||||
return reply, nil
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user