Merge pull request 'A module is a repository and a path, the builder announces what it made, and the mesh acts on it' (#20) from feat/a-module-is-a-repository-and-a-path into main

This commit was merged in pull request #20.
This commit is contained in:
2026-09-13 11:16:32 +02:00
15 changed files with 867 additions and 7 deletions
+57 -3
View File
@@ -56,6 +56,14 @@ is dialled except the broker.
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
MESH_WORKSPACE where to clone and build (default: a temporary directory)
It also builds one module and stops, which is how a mesh is raised — before there is a
broker to take work from or a registry to publish into:
mesh-builder build <repository> [--path P] [--ref COMMIT] [--registry HOST:PORT]
Without --registry the artifacts stay in this machine's container runtime, named by the
digest of their own configuration. The result is printed as JSON.
`
func run() error {
@@ -64,6 +72,8 @@ func run() error {
case "version":
fmt.Println(version)
return nil
case "build":
return buildOnce(context.Background(), os.Args[2:])
default:
fmt.Print(usage)
return nil
@@ -82,9 +92,18 @@ func run() error {
if workspace == "" {
workspace = os.TempDir() + "/mesh-builder"
}
on, err := os.Hostname()
if err != nil {
on = "a build machine"
// **The mesh's name for this machine, not the container's.** A build is reported to the rest
// of the mesh, and a report whose origin reads `104cb10e105b` names something no other module
// can look up. The mesh already knows the answer and has a way to say it — `${machine:name}`
// in the environment file this module is handed — so the hostname is only what is left when
// nobody said.
on := os.Getenv("MESH_NODE")
if on == "" {
hostname, err := os.Hostname()
if err != nil {
return fmt.Errorf("this build machine has no name: nothing said MESH_NODE and the host would not say either: %w", err)
}
on = hostname
}
ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM)
@@ -183,6 +202,7 @@ func answer(ctx context.Context, channel *amqp.Channel, publisher builder.Publis
Name: made.Name, Kind: made.Kind, Reference: made.Reference,
})
}
result.Against = built.Against
fmt.Printf(" built %s from %s\n", built.Manifest.Module, short(built.Commit))
}
}
@@ -211,11 +231,45 @@ func answer(ctx context.Context, channel *amqp.Channel, publisher builder.Publis
}); err != nil {
fmt.Fprintf(os.Stderr, "cannot answer a build request: %v\n", err)
}
// **And announced, which is a different act from answering.** The reply goes to whoever asked
// and is correlated to their request; this says to the whole mesh that a module now exists at
// a commit, and the catalogue places it in the module graph (novox/hq ADR 0072). A build
// nobody asked for still has to be announced, or the graph knows less than the registry does.
//
// Only on success: a failed build produced no module-version, and announcing one would put
// something in the graph that was never made.
if result.Failed == "" && result.Commit != "" {
announced := 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,
}
if err := link.EmitEvent(publishCtx, channel, link.KeyModuleBuilt, "builder", on, announced); err != nil {
// Said, not fatal: the build happened and was answered. A module the catalogue has not
// heard of is a gap somebody can close; a build reported as failed because announcing
// it failed is a lie about work that was done.
fmt.Fprintf(os.Stderr, " built, but could not announce it: %v\n", err)
}
}
// Acknowledged only once the answer is away, so a builder that dies before answering leaves
// the request for another machine rather than losing it.
_ = delivery.Ack(false)
}
// moduleOf reads the module's name out of the manifest it just built, which is the only place it is
// authoritative — the request named a repository and a path, not a module.
func moduleOf(manifest json.RawMessage) string {
var named struct {
Module string `json:"module"`
}
if err := json.Unmarshal(manifest, &named); err != nil {
return ""
}
return named.Module
}
func short(commit string) string {
if len(commit) > 8 {
return commit[:8]
+134
View File
@@ -0,0 +1,134 @@
package main
import (
"context"
"encoding/json"
"errors"
"flag"
"fmt"
"os"
"github.com/novox/mesh-control/internal/builder"
)
// buildOnce is the builder doing one build and stopping, with no broker and no mesh.
//
// **This is how a mesh is raised** (novox/hq ADR 0073). The installer carries this program and runs
// it once, before anything else exists, to produce the control plane from the same repository and
// path that every later rebuild of it will use. What raises the mesh is therefore the same thing
// that maintains it — not a second mechanism that has to be kept in step with the first and is
// exercised once per new mesh, which is how often enough to rot.
//
// It takes no work from a queue and answers nobody: there is no broker yet, and the only thing
// waiting for the answer is the installer that started it. So the result goes to standard output as
// JSON, which is what the installer reads.
//
// With no --registry it publishes nowhere and the image stays in this machine's container runtime,
// named by the digest of its own configuration (see builder.Local). That is the genesis case. With
// one, it publishes as it always does — the same command is how an operator builds a module by hand
// on a mesh that already exists.
func buildOnce(ctx context.Context, args []string) error {
set := flag.NewFlagSet("build", flag.ContinueOnError)
path := set.String("path", "", "the module's directory inside the repository (default: its root)")
ref := set.String("ref", "", "the commit to build; a branch is a moving target somebody else controls")
registry := set.String("registry", "",
"host:port to publish to. Without it the artifacts stay in this machine's container runtime, which is the genesis case")
workspace := set.String("workspace", "", "where to clone and build (default: a temporary directory)")
positionals, err := parseAround(set, args)
if err != nil {
return err
}
if len(positionals) != 1 {
return errors.New("mesh-builder build <repository> [--path P] [--ref COMMIT] [--registry HOST:PORT]")
}
repository := positionals[0]
where := *workspace
if where == "" {
where = os.TempDir() + "/mesh-builder-once"
}
// Local unless told otherwise, because the moment this exists for has nowhere to publish. A
// default pointing at a registry would mean genesis failing at a push to something that is not
// there yet, one step away from the thing that could explain it.
var publisher builder.Publisher = builder.Local{Run: builder.Command}
if *registry != "" {
publisher = builder.Registry{Address: *registry, Run: builder.Command}
}
fmt.Fprintf(os.Stderr, "building %s", repository)
if *path != "" {
fmt.Fprintf(os.Stderr, " at %s", *path)
}
if *ref != "" {
fmt.Fprintf(os.Stderr, " at %s", *ref)
}
fmt.Fprintln(os.Stderr)
built, buildErr := builder.Build(ctx, builder.Command, publisher, repository, *path, *ref, where)
if buildErr != nil {
return buildErr
}
// To standard output, and everything else to standard error, so the caller can read this
// without having to separate it from progress.
out := onceResult{
Module: built.Manifest.Module,
Commit: built.Commit,
Repository: repository,
Path: *path,
Ref: *ref,
Manifest: built.Manifest,
Against: built.Against,
}
for _, made := range built.Built {
out.Made = append(out.Made, madeArtifact{Name: made.Name, Kind: made.Kind, Reference: made.Reference})
}
body, err := json.MarshalIndent(out, "", " ")
if err != nil {
return err
}
fmt.Println(string(body))
fmt.Fprintf(os.Stderr, " built %s from %s\n", built.Manifest.Module, short(built.Commit))
return nil
}
// onceResult is what one build reports to whoever started it.
//
// The same fields the mesh records for a build, so a reader comparing a genesis build against an
// ordinary one is comparing the same thing said the same way.
type onceResult struct {
Module string `json:"module"`
Commit string `json:"commit"`
Repository string `json:"repository"`
Path string `json:"path,omitempty"`
Ref string `json:"ref,omitempty"`
Manifest any `json:"manifest"`
Made []madeArtifact `json:"made"`
Against []string `json:"against,omitempty"`
}
type madeArtifact struct {
Name string `json:"name"`
Kind string `json:"kind"`
Reference string `json:"reference"`
}
// parseAround lets flags appear on either side of the repository, because a person writing this by
// hand will put them wherever reads best and the standard parser stops at the first thing that is
// not a flag. The same helper the control plane's commands use, for the same reason.
func parseAround(set *flag.FlagSet, args []string) ([]string, error) {
var positionals []string
rest := args
for {
if err := set.Parse(rest); err != nil {
return nil, err
}
rest = set.Args()
if len(rest) == 0 {
return positionals, nil
}
positionals = append(positionals, rest[0])
rest = rest[1:]
}
}
+5
View File
@@ -86,6 +86,8 @@ func run() error {
return brokerCommand(args[1:])
case "serve":
return serve(ctx)
case "upgrade":
return upgradeCommand(ctx, args[1:])
case "declare":
return declare(ctx, args[1:])
case "overlay":
@@ -140,6 +142,9 @@ func usage() {
module forget <name> remove one, unless a node runs it or the mesh holds things for it
module forget <name> --and-what-it-holds ...and discard its settings, secrets and ports too
module issue <name> --node <m> a broker account for a module, scoped to its emits and consumes
upgrade <name> what happens when this module's current version moves
upgrade <name> roll-out [--together] ...send it to the machines running it
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
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
+6
View File
@@ -86,6 +86,12 @@ func serve(ctx context.Context) error {
// 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.
server.Records(builds{inv})
// And what the catalogue decided a build meant. The builder's own result is already handled
// above; this is the other half — the control plane is the only one of the three that knows
// which machines run the thing, so it is the one that acts (novox/hq ADR 0072).
if err := server.Follows(following{open}); err != nil {
return err
}
return server.Serve(ctx)
}
+165
View File
@@ -0,0 +1,165 @@
package main
import (
"context"
"errors"
"flag"
"fmt"
"strings"
"github.com/novox/mesh-control/internal/inventory"
"github.com/novox/mesh-control/internal/link"
)
// following acts on what the catalogue announces.
//
// **It holds the stores, not a copy of the decision.** What to do about an upgrade is read when
// one arrives, so changing it takes effect on the next upgrade rather than on the next restart of
// the control plane.
type following struct{ open *stores }
// Upgraded sends the machines running a module the version the catalogue now considers current —
// or records that they are behind, which is the default and needs no record.
//
// **Recording is not a second code path.** A machine that is not running what the mesh would send
// it is already something the mesh notices and reports; that is what `status` and `push --behind`
// are built on. So "record it" is the absence of an action, and the only thing this has to decide
// is whether to act.
func (f following) Upgraded(ctx context.Context, u link.Upgraded) error {
inv := f.open.inventory
decision, err := inv.UpgradeOf(ctx, u.Module)
if err != nil {
return err
}
on, err := inv.Running(ctx, u.Module)
if err != nil {
return err
}
if len(on) == 0 {
fmt.Printf("%s moved to %s; no machine runs it\n", u.Module, shortCommit(u.Commit))
return nil
}
if !decision.RollOut {
// Named rather than counted, and said even though nothing happens: an upgrade that was
// deliberately not rolled out and an upgrade that was never noticed look identical in a
// log that only speaks when it acts.
fmt.Printf("%s moved to %s; %s %s behind it, and this mesh records upgrades rather than "+
"rolling them out — `push --behind` when you want them\n",
u.Module, shortCommit(u.Commit), readableList(on), isAre(len(on)))
return nil
}
if decision.Together {
fmt.Printf("%s moved to %s; sending %s together\n",
u.Module, shortCommit(u.Commit), readableList(on))
return sendTo(ctx, f.open, on)
}
// One at a time, and stopping at the first that fails.
//
// **Stopping is the point.** The machines are done one after another precisely so that a
// version that breaks the first one does not reach the rest; carrying on past a failure would
// make this the same as sending them together, only slower.
fmt.Printf("%s moved to %s; sending %s one at a time\n",
u.Module, shortCommit(u.Commit), readableList(on))
for _, node := range on {
if err := sendTo(ctx, f.open, []string{node}); err != nil {
return fmt.Errorf("%s did not take %s, so the machines after it were left alone: %w",
node, u.Module, err)
}
}
return nil
}
// readableList names machines the way a sentence does, because this is read by a person deciding
// whether an upgrade went where they expected.
func readableList(names []string) string {
switch len(names) {
case 0:
return "nothing"
case 1:
return names[0]
case 2:
return names[0] + " and " + names[1]
}
return strings.Join(names[:len(names)-1], ", ") + " and " + names[len(names)-1]
}
func shortCommit(commit string) string {
if len(commit) > 8 {
return commit[:8]
}
return commit
}
func isAre(n int) string {
if n == 1 {
return "is"
}
return "are"
}
// upgradeCommand says what should happen when a module's current version moves.
func upgradeCommand(ctx context.Context, args []string) error {
set := flag.NewFlagSet("upgrade", flag.ContinueOnError)
together := set.Bool("together", false,
"send every machine running it at once, instead of one after another")
positionals, err := parseAround(set, args)
if err != nil {
return err
}
if len(positionals) == 0 {
return errors.New("upgrade <module> [roll-out|record] [--together]")
}
module := positionals[0]
open, err := openStores(ctx)
if err != nil {
return err
}
defer open.Close()
inv := open.inventory
if len(positionals) == 1 {
decision, err := inv.UpgradeOf(ctx, module)
if err != nil {
return err
}
fmt.Println(sayUpgrade(module, decision))
return nil
}
var decision inventory.Upgrade
switch positionals[1] {
case "roll-out":
decision = inventory.Upgrade{RollOut: true, Together: *together}
case "record":
if *together {
// Refused rather than ignored: --together only means anything for a roll-out, and
// accepting it here would store a preference that never applies and looks like it does.
return errors.New("`--together` says how to roll out, so it cannot be given with " +
"`record`, which is the choice not to")
}
decision = inventory.Upgrade{}
default:
return fmt.Errorf("upgrade <module> roll-out|record — not %q", positionals[1])
}
if err := inv.SetUpgradeOf(ctx, module, decision); err != nil {
return err
}
fmt.Println(sayUpgrade(module, decision))
return nil
}
func sayUpgrade(module string, u inventory.Upgrade) string {
if !u.RollOut {
return fmt.Sprintf("when %s moves, the mesh records it and the machines running it are "+
"reported as behind", module)
}
if u.Together {
return fmt.Sprintf("when %s moves, every machine running it is sent the new version "+
"together", module)
}
return fmt.Sprintf("when %s moves, the machines running it are sent the new version one at "+
"a time, stopping at the first that fails", module)
}
+6 -3
View File
@@ -242,11 +242,14 @@ func (m *Management) CreateBuilderAccount(ctx context.Context, name, password st
// 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 + "$",
// The exchange, and nothing else. **Not the default exchange**: permission there is
// 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 exchange, which is why they can.
"write": "^" + regexp.QuoteMeta(ExchangeName) + "$",
// 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 {
+45 -1
View File
@@ -11,6 +11,7 @@ import (
"os"
"os/exec"
"path/filepath"
"regexp"
"sort"
"strings"
@@ -42,6 +43,15 @@ type Publisher interface {
// Result is everything one build produced.
type Result struct {
// Against is every pinned image this build was built on top of, read out of its own inputs.
//
// **Derived, not declared** (novox/hq ADR 0009): a declared list of dependencies drifts from
// what the code actually uses, and an artifact is out of date when anything it was built
// against moved. These are artifact references rather than module-versions, because that is
// what a build input names; resolving them to modules is the catalogue's work, since it is
// what knows which module-version published which artifact.
Against []string
// Manifest is the module as the mesh should hold it: artifacts resolved to digests.
Manifest catalogue.Manifest
// Commit is what was built, so "is this current?" is answerable without building again.
@@ -124,7 +134,8 @@ func Build(ctx context.Context, run Runner, publish Publisher,
if err != nil {
return Result{}, err
}
return Result{Manifest: resolved, Commit: commit, Built: built}, nil
return Result{Manifest: resolved, Commit: commit, Built: built,
Against: against(within, manifest)}, nil
}
// inside resolves a module's path within a clone, and refuses one that leaves it.
@@ -157,6 +168,39 @@ func describe(path string) string {
return path
}
// pinnedImage matches an image reference pinned by digest, which is the only kind a build input is
// allowed to name — a tag is something somebody else can move under you.
var pinnedImage = regexp.MustCompile(`[A-Za-z0-9][A-Za-z0-9._/:-]*@sha256:[0-9a-f]{64}`)
// against reads what this module's image artifacts are built on top of, out of the files that
// build them. Nothing is guessed: a reference that is not written down is not reported.
func against(within string, manifest catalogue.Manifest) []string {
if manifest.Build == nil {
return nil
}
seen := map[string]bool{}
var out []string
for _, a := range manifest.Build.Artifacts {
if a.Kind != catalogue.ArtifactImage || a.From == "" {
continue
}
body, err := os.ReadFile(filepath.Join(within, a.From))
if err != nil {
// Not fatal: the build itself already failed if this file was needed and missing, and
// reporting no edges is honest where inventing them would not be.
continue
}
for _, found := range pinnedImage.FindAllString(string(body), -1) {
if !seen[found] {
seen[found] = true
out = append(out, found)
}
}
}
sort.Strings(out)
return out
}
// ManifestName is the one file a module repository must have.
//
// At the root, and named the same in every repository. A convention somebody can look for beats a
+61
View File
@@ -0,0 +1,61 @@
package builder
import (
"context"
"fmt"
"strings"
)
// Local is what a build publishes into when there is nowhere to publish yet.
//
// **This exists for exactly one moment: raising a mesh** (novox/hq ADR 0073). The installer builds
// the control plane on the machine that is about to run it, and at that moment there is no registry
// — the registry is installed afterwards, by the control plane this build produces. So the artifact
// stays where the build left it: in the machine's own container runtime.
//
// That is not a weaker kind of pinning. An image held locally is named by the digest of its own
// configuration, which is content-addressed and unforgeable and requires nothing to have served it
// — the same identity the installer has always used for the image it carried. What changes when a
// registry exists is not that the artifact becomes exact, but that something other than this
// machine can fetch it.
//
// It is deliberately unable to publish an archive. An archive has no local identity to fall back
// on: it is bytes that only mean something once something serves them at a URL. A build that
// produces one before there is anywhere to put it has produced nothing usable, and saying so is
// better than returning a path on a disk that no other machine can read.
type Local struct {
// Run is how docker is invoked, so a test does not need one.
Run Runner
}
// PublishImage leaves the image where the build put it, and names it by its own configuration.
//
// The local tag is not returned: a tag is a name somebody can move, and every other reference in a
// resolved manifest is exact. The digest is read back from the runtime rather than computed, for
// the same reason the registry publisher reads it back from the registry — what matters is what
// will be served for that reference, and only the thing serving it can say.
func (l Local) PublishImage(ctx context.Context, localTag, repository string) (string, error) {
out, err := l.Run(ctx, "", "docker", "image", "inspect", "--format", "{{.Id}}", localTag)
if err != nil {
return "", fmt.Errorf("cannot read back the image just built as %s: %w", localTag, err)
}
id := strings.TrimSpace(out)
if !strings.HasPrefix(id, "sha256:") || len(id) != len("sha256:")+64 {
// Refused rather than passed on. A bundle naming an image by anything a person could move
// is refused by the machine applying it, and a value that is not an identity would fail
// there instead — one step further from the thing that could explain it.
return "", fmt.Errorf(
"the container runtime named the image just built %q, which is not an image id: "+
"sha256 and sixty-four hex characters", id)
}
return id, nil
}
// PublishArchive refuses, and says why rather than inventing somewhere to put it.
func (l Local) PublishArchive(ctx context.Context, repository string, body []byte, digest string) (string, error) {
return "", fmt.Errorf(
"%s declares an archive, and this build has nowhere to publish one. An image can stay in "+
"the machine's own runtime and still be named exactly; an archive is bytes that mean "+
"nothing until something serves them. Build this once the mesh has a registry",
repository)
}
+29
View File
@@ -456,6 +456,11 @@ func (r Resolution) Declaration(with Rendering) ([]map[string]any, error) {
if err := pinned(copied, m.Module); err != nil {
return nil, err
}
// And the other half of the same question: an image nobody has published yet is named
// by an artifact rather than by a placeholder digest, and is just as unrunnable.
if err := built(copied, m.Module); err != nil {
return nil, err
}
publishedOn(copied, m.Module, with)
copied["id"] = m.Module + "." + fmt.Sprint(resource["id"])
// A service saying what it reflects names resources within its own module, so those
@@ -887,6 +892,30 @@ func pinned(resource map[string]any, module string) error {
module, resource["id"])
}
// built refuses a container that still names an artifact nobody has made.
//
// **Said here, by the thing that knows what "artifact" means.** A container naming an artifact is
// a module saying "the mesh builds this"; the field is resolved into an image when a build
// publishes one, and until then there is nothing to run. A host receiving it refuses the whole
// declaration — correctly, since its language has no such field — but what it can say is that a
// container does not use "artifact", which tells a reader the manifest is malformed. It is not:
// it is unbuilt, which is a different problem with a different fix, and only the mesh is in a
// position to tell them apart (novox/hq ADR 0073).
func built(resource map[string]any, module string) error {
if fmt.Sprint(resource["type"]) != "container" {
return nil
}
artifact, ok := resource["artifact"].(string)
if !ok || artifact == "" {
return nil
}
return fmt.Errorf(
"%s has not been built for this mesh: %v is its %q artifact, and no build has published "+
"one. Build it — `build <repository> --path <path>` — and the module will name what "+
"came out instead",
module, resource["id"], artifact)
}
// publishedOn puts the machine's own port on the outside of a container's mapping.
//
// **A module writes the port the software uses; the mesh says where the machine puts it**
+75
View File
@@ -796,3 +796,78 @@ func (i *Inventory) Catalogued(ctx context.Context) ([]Entry, error) {
// providedBy is what the source column says for a module the control plane ships.
const providedBy = "the control plane"
// Upgrade is what the mesh decided to do when a module's current version moves.
type Upgrade struct {
// RollOut is true when the machines running it should be sent the new version. False means
// record it and stop — which needs no record of its own, because a machine not running what
// the mesh would send it is already something the mesh reports.
RollOut bool
// Together is true when every machine running it is sent the new version at once, rather than
// one after another. Only meaningful when RollOut is.
Together bool
}
// UpgradeOf is what to do when this module moves.
//
// A module the mesh does not hold is not an error here: the catalogue may know of modules this
// mesh has never registered, and being told one of them moved is information, not a fault. The
// answer is the safe one — record it — because there is nothing to roll out to.
func (i *Inventory) UpgradeOf(ctx context.Context, module string) (Upgrade, error) {
var u Upgrade
var policy string
err := i.store.Pool().QueryRow(ctx,
`select upgrade, upgrade_together from module where name = $1`, module).
Scan(&policy, &u.Together)
if errors.Is(err, pgx.ErrNoRows) {
return Upgrade{}, nil
}
if err != nil {
return Upgrade{}, err
}
u.RollOut = policy == "roll-out"
return u, nil
}
// SetUpgradeOf records what to do when this module moves.
func (i *Inventory) SetUpgradeOf(ctx context.Context, module string, u Upgrade) error {
policy := "record"
if u.RollOut {
policy = "roll-out"
}
tag, err := i.store.Pool().Exec(ctx,
`update module set upgrade = $2, upgrade_together = $3 where name = $1`,
module, policy, u.Together)
if err != nil {
return err
}
if tag.RowsAffected() == 0 {
return fmt.Errorf("this mesh holds no module called %s", module)
}
return nil
}
// Running is every machine assigned a module, in a stable order.
//
// **Assigned, not reported.** A machine that is assigned the module and has not applied it yet is
// exactly the machine an upgrade most needs to reach; waiting for it to report the old version
// first would mean the machines furthest behind are the last to be caught up.
func (i *Inventory) Running(ctx context.Context, module string) ([]string, error) {
rows, err := i.store.Pool().Query(ctx,
`select n.name from assignment a join node n on n.id = a.node
where a.module = $1 order by n.name`, module)
if err != nil {
return nil, err
}
defer rows.Close()
var out []string
for rows.Next() {
var name string
if err := rows.Scan(&name); err != nil {
return nil, err
}
out = append(out, name)
}
return out, rows.Err()
}
@@ -0,0 +1,26 @@
-- What the mesh should do when the catalogue says a module has been upgraded.
--
-- novox/hq ADR 0072. The builder announces a build, the catalogue decides whether that is an
-- upgrade, and the control plane is the only one of the three that knows which machines are
-- running the thing. So it is the one that acts -- and what it does has to be a decision somebody
-- made, not a behaviour compiled in.
--
-- Two separate questions, deliberately not one:
--
-- whether -- roll the new version out, or record that the machine is behind and stop there.
-- how -- one machine at a time, or all of them together.
--
-- They are separate because the safe answer to the first is not the safe answer to the second: a
-- mesh may well want every upgrade applied automatically and still never want its only two
-- machines restarted in the same breath.
--
-- Defaulted to recording rather than rolling out, and per module rather than mesh-wide. A mesh
-- that upgrades everything it builds the moment it builds it is a reasonable thing to want and a
-- terrible thing to arrive by default -- the first module to inherit it would be the control plane
-- itself, upgrading itself out from under the push that was applying it.
alter table module add column upgrade text not null default 'record'
check (upgrade in ('record', 'roll-out'));
-- Only meaningful when upgrade is 'roll-out'. Kept anyway when it is not, so turning roll-out on
-- does not silently also decide this.
alter table module add column upgrade_together boolean not null default false;
+4
View File
@@ -82,6 +82,10 @@ type BuildResult struct {
// Made is each artifact, for reporting.
Made []MadeArtifact `json:"made,omitempty"`
// Against is every pinned image this was built on top of, read out of the build's own inputs
// (novox/hq ADR 0009). The catalogue turns these into edges; nothing else need care.
Against []string `json:"against,omitempty"`
// Failed is why, when it did.
Failed string `json:"failed,omitempty"`
}
+92
View File
@@ -0,0 +1,92 @@
package link
import (
"context"
"crypto/rand"
"encoding/hex"
"encoding/json"
"fmt"
"time"
amqp "github.com/rabbitmq/amqp091-go"
)
// Emitting a module event from Go.
//
// **Every event rides one topic exchange** (novox/hq ADR 0042), which is not the direct exchange
// nodes and the control plane speak over. A module that announces something publishes here, and
// consumers bind their own durable queue to a pattern over it.
//
// This exists because the builder is a module written in Go while every other emitter is
// TypeScript on the sdk. The envelope is the sdk's, reproduced exactly: the body is the payload
// alone and everything about the event travels as headers. A second shape would be a second thing
// for consumers to handle, and they are written against the first.
const (
// EventsExchange is where every event rides. Named here rather than imported from the broker
// package for the same reason BuildQueueName is duplicated there — one direction of dependency.
EventsExchange = "mesh.events"
)
// EmitEvent publishes one module event, in the envelope the sdk's consumers expect.
//
// Persistent, because an event that a broker restart loses is not an announcement. The publish is
// not confirmed here: the caller has already done the work the event describes, and a build that
// succeeded must not be reported as failed because saying so failed.
func EmitEvent(ctx context.Context, channel *amqp.Channel, eventType, source, node string, body any) error {
payload, err := json.Marshal(body)
if err != nil {
return fmt.Errorf("cannot serialise a %s event: %w", eventType, err)
}
id, err := eventID()
if err != nil {
return err
}
return channel.PublishWithContext(ctx, EventsExchange, eventType, false, false, amqp.Publishing{
ContentType: "application/json",
DeliveryMode: amqp.Persistent,
MessageId: id,
Timestamp: time.Now().UTC(),
Body: payload,
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",
},
})
}
// eventID is what a consumer deduplicates on: delivery is at-least-once, so a handler must be able
// to tell a redelivery from a second event, and only the emitter can say which it is.
func eventID() (string, error) {
raw := make([]byte, 16)
if _, err := rand.Read(raw); err != nil {
return "", fmt.Errorf("cannot make an event id: %w", err)
}
return hex.EncodeToString(raw), nil
}
// KeyModuleBuilt is what the builder announces when it has built something. The catalogue places
// it in the module graph; nothing else need care.
const KeyModuleBuilt = "module.builder.built"
// KeyModuleUpgraded is the catalogue saying a module's current version has moved.
//
// **The control plane hooks the meaning, not the build.** The builder says what it built; the
// catalogue decides whether that was an upgrade — a rebuild producing the commit already current
// is not one — and only this says anything the control plane can act on. Consuming the build
// directly would make the control plane re-derive a decision another module already made, and the
// two would eventually disagree (novox/hq ADR 0072).
const KeyModuleUpgraded = "module.mesh-catalog.upgraded"
// UpgradeQueue is where those land. Durable and named, not a temporary queue: an upgrade announced
// while the control plane is restarting is exactly the one that must not be missed.
const UpgradeQueue = "control.upgrades"
// Upgraded is what the catalogue says when a module's current version moves.
type Upgraded struct {
Module string `json:"module"`
Commit string `json:"commit"`
Previous string `json:"previous"`
}
+94
View File
@@ -50,6 +50,7 @@ type Server struct {
listener Listener
recorder Recorder
log *log.Logger
upgrader Upgrader
}
// Records tells the server where to keep build results.
@@ -59,6 +60,35 @@ type Server struct {
// making it supply one would have it construct something it never uses.
func (s *Server) Records(r Recorder) { s.recorder = r }
// Upgrader is what the control plane does when the catalogue says a module moved.
//
// An interface for the same reason Enroller is one: deciding what an upgrade means for the
// machines running it is a different concern from noticing that one was announced, and only the
// first needs a database.
type Upgrader interface {
// Upgraded is told which module moved and between which commits. An error is logged and the
// message is not requeued: an upgrade the control plane could not act on is not one it will
// act on by being handed the same message again, and a poison message on a durable queue
// would stop every upgrade behind it.
Upgraded(ctx context.Context, u Upgraded) error
}
// Follows says what to do about upgrades, and binds the queue they arrive on.
//
// **Not bound unless something is listening.** A durable queue bound to every upgrade with no
// consumer fills up quietly, and the first symptom is a broker out of disk rather than anything
// about modules.
func (s *Server) Follows(u Upgrader) error {
if _, err := s.channel.QueueDeclare(UpgradeQueue, true, false, false, false, nil); err != nil {
return fmt.Errorf("cannot declare the %s queue: %w", UpgradeQueue, err)
}
if err := s.channel.QueueBind(UpgradeQueue, KeyModuleUpgraded, EventsExchange, false, nil); err != nil {
return fmt.Errorf("cannot bind %s to %s/%s: %w", UpgradeQueue, EventsExchange, KeyModuleUpgraded, err)
}
s.upgrader = u
return nil
}
// Connect opens the control plane's own connection to the broker.
func Connect(enroller Enroller, listener Listener) (*Server, error) {
url := strings.TrimSpace(os.Getenv(AMQPVar))
@@ -88,6 +118,13 @@ func Connect(enroller Enroller, listener Listener) (*Server, error) {
conn.Close()
return nil, fmt.Errorf("cannot declare the %s exchange: %w", Exchange, err)
}
// The events exchange too. The control plane 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)
@@ -138,14 +175,39 @@ func (s *Server) Serve(ctx context.Context) error {
return err
}
// The upgrade queue, when something is listening for them. A second queue rather than a
// second consumer on the first: two consumers on one queue split its messages between them,
// which is the fault the comment above exists about. Two queues share nothing.
var upgrades <-chan amqp.Delivery
if s.upgrader != nil {
upgrades, err = s.channel.ConsumeWithContext(ctx, UpgradeQueue, "control-plane-upgrades",
false, false, false, false, nil)
if err != nil {
return err
}
}
closed := s.conn.NotifyClose(make(chan *amqp.Error, 1))
s.log.Printf("consuming %s, bound to %s/{%s,%s,%s,%s}",
ControlQueue, Exchange, KeyEnrol, KeyReport, KeyAlive, KeyBuilt)
if s.upgrader != nil {
s.log.Printf("consuming %s, bound to %s/%s", UpgradeQueue, EventsExchange, KeyModuleUpgraded)
}
for {
select {
case <-ctx.Done():
return nil
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 broker stopped.
if !ok {
if upgrades != nil {
return errors.New("the broker stopped delivering upgrades")
}
continue
}
s.upgraded(ctx, delivery)
case reason := <-closed:
// Said rather than returned quietly. A control plane whose broker connection dropped
// is a mesh where nothing can be told anything, and the reason is the first thing
@@ -309,3 +371,35 @@ func (s *Server) handleBuilt(ctx context.Context, delivery amqp.Delivery) {
}
_ = delivery.Ack(false)
}
// upgraded hands one announcement to whatever is following them.
//
// **Acknowledged whatever happens.** A failure here is the control plane being unable to act on an
// upgrade — a machine that cannot be resolved, a broker that will not take a declaration — and
// none of those get better by being handed the same message again. Requeuing would put a poison
// message at the head of a durable queue and stop every upgrade behind it, which turns one module
// nobody can push into a mesh that stops following its own catalogue.
func (s *Server) upgraded(ctx context.Context, delivery amqp.Delivery) {
defer func() { _ = delivery.Ack(false) }()
var u Upgraded
if err := json.Unmarshal(delivery.Body, &u); err != nil {
s.log.Printf("an upgrade announcement could not be read: %v", err)
return
}
if u.Module == "" {
s.log.Printf("an upgrade announcement named no module; ignored")
return
}
if err := s.upgrader.Upgraded(ctx, u); err != nil {
s.log.Printf("%s moved to %s and the mesh could not act on it: %v",
u.Module, short(u.Commit), err)
}
}
// short is a commit as people read it.
func short(commit string) string {
if len(commit) > 8 {
return commit[:8]
}
return commit
}
+68
View File
@@ -0,0 +1,68 @@
{
"module": "mesh-control",
"version": "1",
"slug": "control",
"capabilities": [
"container-runtime"
],
"claims": [
{
"name": "the-control-plane",
"scope": "mesh"
}
],
"own-secrets": {
"inventory": "/var/lib/mesh/mesh-control/inventory",
"identity": "/var/lib/mesh/mesh-control/identity",
"licences": "/var/lib/mesh/mesh-control/licences",
"broker": "/var/lib/mesh/mesh-control/broker",
"broker-management": "/var/lib/mesh/mesh-control/broker-management",
"broker-address": "/var/lib/mesh/mesh-control/broker-address"
},
"resources": [
{
"id": "mesh-state",
"type": "directory",
"path": "/var/lib/mesh/mesh-control",
"mode": "0700"
},
{
"id": "control-env",
"type": "file",
"path": "/var/lib/mesh/mesh-control/control.env",
"mode": "0600",
"content": "MESH_STORE_INVENTORY=${secret:inventory}\nMESH_STORE_IDENTITY=${secret:identity}\nMESH_STORE_LICENCES=${secret:licences}\nMESH_BROKER_AMQP=${secret:broker}\nMESH_BROKER_MANAGEMENT=${secret:broker-management}\nMESH_BROKER_ADDRESS=${secret:broker-address}\n"
},
{
"id": "server",
"type": "container",
"name": "mesh-control",
"network": "host",
"args": [
"serve"
],
"env-file": [
"/var/lib/mesh/mesh-control/control.env"
],
"env": {
"MESH_BROKER_CERTIFICATE": "/broker-tls/tls.crt"
},
"volumes": [
"mesh-broker-tls:/broker-tls:ro"
],
"restart-on": [
"control-env"
],
"artifact": "server"
}
],
"build": {
"artifacts": [
{
"name": "server",
"kind": "image",
"from": "Dockerfile"
}
]
}
}