A module is a repository and a path, the builder announces what it made, and the mesh acts on it #20
@@ -56,6 +56,14 @@ is dialled except the broker.
|
|||||||
MESH_REGISTRY host:port to publish artifacts to, when the mesh has not said
|
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_BINDING a file the mesh wrote saying where the artifact store is
|
||||||
MESH_WORKSPACE where to clone and build (default: a temporary directory)
|
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 {
|
func run() error {
|
||||||
@@ -64,6 +72,8 @@ func run() error {
|
|||||||
case "version":
|
case "version":
|
||||||
fmt.Println(version)
|
fmt.Println(version)
|
||||||
return nil
|
return nil
|
||||||
|
case "build":
|
||||||
|
return buildOnce(context.Background(), os.Args[2:])
|
||||||
default:
|
default:
|
||||||
fmt.Print(usage)
|
fmt.Print(usage)
|
||||||
return nil
|
return nil
|
||||||
@@ -82,9 +92,18 @@ func run() error {
|
|||||||
if workspace == "" {
|
if workspace == "" {
|
||||||
workspace = os.TempDir() + "/mesh-builder"
|
workspace = os.TempDir() + "/mesh-builder"
|
||||||
}
|
}
|
||||||
on, err := os.Hostname()
|
// **The mesh's name for this machine, not the container's.** A build is reported to the rest
|
||||||
if err != nil {
|
// of the mesh, and a report whose origin reads `104cb10e105b` names something no other module
|
||||||
on = "a build machine"
|
// 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)
|
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,
|
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))
|
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 {
|
}); err != nil {
|
||||||
fmt.Fprintf(os.Stderr, "cannot answer a build request: %v\n", err)
|
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
|
// Acknowledged only once the answer is away, so a builder that dies before answering leaves
|
||||||
// the request for another machine rather than losing it.
|
// the request for another machine rather than losing it.
|
||||||
_ = delivery.Ack(false)
|
_ = 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 {
|
func short(commit string) string {
|
||||||
if len(commit) > 8 {
|
if len(commit) > 8 {
|
||||||
return commit[:8]
|
return commit[:8]
|
||||||
|
|||||||
@@ -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:]
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -86,6 +86,8 @@ func run() error {
|
|||||||
return brokerCommand(args[1:])
|
return brokerCommand(args[1:])
|
||||||
case "serve":
|
case "serve":
|
||||||
return serve(ctx)
|
return serve(ctx)
|
||||||
|
case "upgrade":
|
||||||
|
return upgradeCommand(ctx, args[1:])
|
||||||
case "declare":
|
case "declare":
|
||||||
return declare(ctx, args[1:])
|
return declare(ctx, args[1:])
|
||||||
case "overlay":
|
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> 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 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
|
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
|
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
|
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
|
||||||
|
|||||||
@@ -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`
|
// 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})
|
||||||
|
// 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)
|
return server.Serve(ctx)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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)
|
||||||
|
}
|
||||||
@@ -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
|
// 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.
|
// already run, which makes the order they start in matter.
|
||||||
"configure": "^" + builds + "$",
|
"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
|
// 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
|
// 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.
|
// not have. Answers go through the node exchange; announcing what was built goes through
|
||||||
"write": "^" + regexp.QuoteMeta(ExchangeName) + "$",
|
// 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.
|
// The build queue and nothing else. Not another machine's declarations.
|
||||||
"read": "^" + builds + "$",
|
"read": "^" + builds + "$",
|
||||||
}); err != nil {
|
}); err != nil {
|
||||||
|
|||||||
@@ -11,6 +11,7 @@ import (
|
|||||||
"os"
|
"os"
|
||||||
"os/exec"
|
"os/exec"
|
||||||
"path/filepath"
|
"path/filepath"
|
||||||
|
"regexp"
|
||||||
"sort"
|
"sort"
|
||||||
"strings"
|
"strings"
|
||||||
|
|
||||||
@@ -42,6 +43,15 @@ type Publisher interface {
|
|||||||
|
|
||||||
// Result is everything one build produced.
|
// Result is everything one build produced.
|
||||||
type Result struct {
|
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 is the module as the mesh should hold it: artifacts resolved to digests.
|
||||||
Manifest catalogue.Manifest
|
Manifest catalogue.Manifest
|
||||||
// Commit is what was built, so "is this current?" is answerable without building again.
|
// 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 {
|
if err != nil {
|
||||||
return Result{}, err
|
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.
|
// 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
|
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.
|
// 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
|
// At the root, and named the same in every repository. A convention somebody can look for beats a
|
||||||
|
|||||||
@@ -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)
|
||||||
|
}
|
||||||
@@ -456,6 +456,11 @@ func (r Resolution) Declaration(with Rendering) ([]map[string]any, error) {
|
|||||||
if err := pinned(copied, m.Module); err != nil {
|
if err := pinned(copied, m.Module); err != nil {
|
||||||
return nil, err
|
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)
|
publishedOn(copied, m.Module, with)
|
||||||
copied["id"] = m.Module + "." + fmt.Sprint(resource["id"])
|
copied["id"] = m.Module + "." + fmt.Sprint(resource["id"])
|
||||||
// A service saying what it reflects names resources within its own module, so those
|
// 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"])
|
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.
|
// 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**
|
// **A module writes the port the software uses; the mesh says where the machine puts it**
|
||||||
|
|||||||
@@ -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.
|
// providedBy is what the source column says for a module the control plane ships.
|
||||||
const providedBy = "the control plane"
|
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;
|
||||||
@@ -82,6 +82,10 @@ type BuildResult struct {
|
|||||||
// Made is each artifact, for reporting.
|
// Made is each artifact, for reporting.
|
||||||
Made []MadeArtifact `json:"made,omitempty"`
|
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 is why, when it did.
|
||||||
Failed string `json:"failed,omitempty"`
|
Failed string `json:"failed,omitempty"`
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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"`
|
||||||
|
}
|
||||||
@@ -50,6 +50,7 @@ type Server struct {
|
|||||||
listener Listener
|
listener Listener
|
||||||
recorder Recorder
|
recorder Recorder
|
||||||
log *log.Logger
|
log *log.Logger
|
||||||
|
upgrader Upgrader
|
||||||
}
|
}
|
||||||
|
|
||||||
// Records tells the server where to keep build results.
|
// 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.
|
// making it supply one would have it construct something it never uses.
|
||||||
func (s *Server) Records(r Recorder) { s.recorder = r }
|
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.
|
// Connect opens the control plane's own connection to the broker.
|
||||||
func Connect(enroller Enroller, listener Listener) (*Server, error) {
|
func Connect(enroller Enroller, listener Listener) (*Server, error) {
|
||||||
url := strings.TrimSpace(os.Getenv(AMQPVar))
|
url := strings.TrimSpace(os.Getenv(AMQPVar))
|
||||||
@@ -88,6 +118,13 @@ func Connect(enroller Enroller, listener Listener) (*Server, error) {
|
|||||||
conn.Close()
|
conn.Close()
|
||||||
return nil, fmt.Errorf("cannot declare the %s exchange: %w", Exchange, err)
|
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 {
|
if _, err := channel.QueueDeclare(ControlQueue, true, false, false, false, nil); err != nil {
|
||||||
conn.Close()
|
conn.Close()
|
||||||
return nil, fmt.Errorf("cannot declare the %s queue: %w", ControlQueue, err)
|
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
|
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))
|
closed := s.conn.NotifyClose(make(chan *amqp.Error, 1))
|
||||||
s.log.Printf("consuming %s, bound to %s/{%s,%s,%s,%s}",
|
s.log.Printf("consuming %s, bound to %s/{%s,%s,%s,%s}",
|
||||||
ControlQueue, Exchange, KeyEnrol, KeyReport, KeyAlive, KeyBuilt)
|
ControlQueue, Exchange, KeyEnrol, KeyReport, KeyAlive, KeyBuilt)
|
||||||
|
if s.upgrader != nil {
|
||||||
|
s.log.Printf("consuming %s, bound to %s/%s", UpgradeQueue, EventsExchange, KeyModuleUpgraded)
|
||||||
|
}
|
||||||
|
|
||||||
for {
|
for {
|
||||||
select {
|
select {
|
||||||
case <-ctx.Done():
|
case <-ctx.Done():
|
||||||
return nil
|
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:
|
case reason := <-closed:
|
||||||
// Said rather than returned quietly. A control plane whose broker connection dropped
|
// 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
|
// 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)
|
_ = 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
@@ -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"
|
||||||
|
}
|
||||||
|
]
|
||||||
|
}
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user