Files
mesh-controller/cmd/mesh-controller/merge_gate.go
T
jochen 58ebe590a5 Ask the planner what a pull request reaches; map a changed file onto modules in one place (hq ADR 0238)
touchedBy is now the only mapping of changed files onto modules — touched, added, and
read by no build — and reachOfMerge the planner's whole answer with the dependency walk.
The merge handler, the plan what-if, the merge gate's width and composition, and a pull
request's check all ask it, so planning and gating cannot disagree. The gate composes
the definitions of the modules a merge would rebuild or add, not every one in the tree,
and the check says the dependents a merge would build after them.
2026-10-06 22:34:57 +02:00

1117 lines
39 KiB
Go

package main
import (
"context"
"crypto/ecdh"
"crypto/rand"
"crypto/sha256"
"encoding/base64"
"encoding/json"
"errors"
"flag"
"fmt"
"io"
"io/fs"
"os"
"path/filepath"
"slices"
"sort"
"strings"
"time"
"github.com/jackc/pgx/v5"
"github.com/novox/mesh-controller/internal/artifacts"
"github.com/novox/mesh-controller/internal/catalogue"
snapshot "github.com/novox/mesh-controller/internal/facts"
"github.com/novox/mesh-controller/internal/inventory"
"github.com/novox/mesh-controller/internal/link"
"github.com/novox/mesh-controller/internal/store"
)
// The merge gate (novox/hq to-be 45 §9, ADR 0227 rule 9): **every machine of the mesh composed with the
// change, and the node-engine's own validator run over each, before the change merges.**
//
// A change is judged against the mesh as the facts snapshot says it runs, twice: once as it is (the
// base) and once with the change applied, each in a throwaway store raised from the snapshot through
// the controller's own code — the same registration, assignment, seats, pins and settings a running
// controller keeps, and the same composition a push makes. What the change breaks is what composes or
// validates in the base and not with the change, named by the machine's role and the module; what was
// already broken is said and does not fail the change.
//
// What it catches, by the incidents it was written from:
// - a manifest the node-engine refuses whole, found only by assigning it to a live machine (236);
// - an identity a real machine's name makes too long for a provision, which refused the provider's
// whole machine (263) — or, since ADR 0225, leaves the consumer out of its grants: either fails;
// - a manifest the controller the mesh runs cannot read (version skew): a catalogue change is judged
// by that controller, so a field only a newer one knows fails here, not at registration;
// - a module removed from its source while machines still run it (ADR 0236), and two definitions of
// one module name (239);
// - a setting the mesh holds that the changed definition cannot compose with (ADR 0163);
// - and, said rather than failed, how wide the rebuild of the merge would be (278).
//
// It runs where a build runs: the build seat, asked to check a pull request (checks.go), so no CI
// outside the mesh is needed; and by hand, with a snapshot and a throwaway store, by anyone.
// mergeCheckInput is what one judging needs.
type mergeCheckInput struct {
facts snapshot.Facts
// repository is the change's owner/repository; tree its checkout, empty when the change is to code
// alone and every module stays as the mesh holds it.
repository, tree string
// changed are the paths the change touches, for the width of its rebuild.
changed []string
// admin is a PostgreSQL the gate may create and drop databases in.
admin string
// say is where progress goes; the verdict is returned.
say io.Writer
}
// mergeVerdict is what the gate found.
type mergeVerdict struct {
Verdict string `json:"verdict"` // pass, warning or fail
Summary string `json:"summary"`
// Judge is the controller build that judged it, and Facts when the snapshot it judged against was taken.
Judge string `json:"judge"`
Facts time.Time `json:"facts"`
Failures []string `json:"failures,omitempty"`
Warnings []string `json:"warnings,omitempty"`
Notes []string `json:"notes,omitempty"`
Machines []mergeMachine `json:"machines"`
Modules map[string]string `json:"modules,omitempty"` // changed module → what changed: added, changed, removed
Width *mergeWidthOf `json:"width,omitempty"`
}
// mergeMachine is one machine, composed as it is and with the change.
type mergeMachine struct {
Name string `json:"name"`
Described string `json:"described"`
Live bool `json:"live-composes"`
Base mergeComposed `json:"base"`
Change mergeComposed `json:"change"`
Added []string `json:"added,omitempty"`
Removed []string `json:"removed,omitempty"`
Modules map[string]int `json:"-"`
}
// mergeComposed is one machine's composition in one store.
type mergeComposed struct {
Composes bool `json:"composes"`
Problems []string `json:"problems,omitempty"`
Withheld []string `json:"withheld,omitempty"`
Unbound []string `json:"unbound,omitempty"`
LeftOut map[string]string `json:"left-out,omitempty"`
Resources []string `json:"resources,omitempty"`
}
// mergeWidthOf is how much a merge of the change would rebuild.
type mergeWidthOf struct {
Modules []string `json:"modules"`
Tiers int `json:"tiers"`
// Unread is the changed paths no module's build reads — in no module's directory — which rebuild
// nothing (issue 280): said, so a reader sees why a change to them moves no module.
Unread []string `json:"unread,omitempty"`
}
// wideRebuild is how many modules a rebuild may take before the gate says so as a warning.
var wideRebuild = 12
// mergeGateCommand is `gate`:
//
// merge-gate --facts <file> --store <postgres> [--repository owner/repo --tree <dir>] [--changed a,b] [--json]
//
// --facts may be a file or `store`, the snapshot the artifact store holds (MESH_REGISTRY names it).
// --store is a PostgreSQL the gate may create and drop databases in: a throwaway, never the mesh's.
func mergeGateCommand(ctx context.Context, args []string) error {
set := flag.NewFlagSet("merge-gate", flag.ContinueOnError)
factsFrom := set.String("facts", "", "the facts snapshot: a file, or `store` for the one the artifact store holds")
admin := set.String("store", os.Getenv("MESH_GATE_POSTGRES"), "a throwaway PostgreSQL the gate may create databases in")
repository := set.String("repository", "", "owner/repository of the change")
tree := set.String("tree", "", "the change's checkout: every module.json in it replaces the mesh's")
changed := set.String("changed", "", "the paths the change touches, comma-separated, for its rebuild's width")
asJSON := set.Bool("json", false, "print the verdict as JSON")
if _, err := parseAround(set, args); err != nil {
return err
}
if *factsFrom == "" || *admin == "" {
return errors.New("merge-gate --facts <file|store> --store <postgres> [--repository owner/repo --tree <dir>] " +
"[--changed paths] [--json]")
}
if (*tree == "") != (*repository == "") && *tree != "" {
return errors.New("--tree names the change's checkout; --repository says which repository it is")
}
f, err := readFacts(ctx, *factsFrom)
if err != nil {
return err
}
out := io.Writer(os.Stdout)
if *asJSON {
out = os.Stderr
}
v, err := judgeChange(ctx, mergeCheckInput{facts: f, repository: *repository, tree: *tree, changed: splitList(*changed),
admin: *admin, say: out})
if err != nil {
return err
}
if *asJSON {
body, err := json.MarshalIndent(v, "", " ")
if err != nil {
return err
}
fmt.Println(string(body))
// And the report a person reads, beside it: what the pull request is told comes from here.
fmt.Fprint(os.Stderr, v.Report())
} else {
fmt.Print(v.Report())
}
if v.Verdict == "fail" {
return errMergeGateFailed
}
return nil
}
// errMergeGateFailed is the gate's own failure: the verdict says why, so the error says nothing more.
var errMergeGateFailed = errors.New("the change fails the merge gate")
// readFacts reads a snapshot from a file, or from the artifact store.
func readFacts(ctx context.Context, from string) (snapshot.Facts, error) {
var body []byte
var err error
if from == "store" {
address := strings.TrimSpace(os.Getenv("MESH_REGISTRY"))
if address == "" {
return snapshot.Facts{}, errors.New("--facts store reads the artifact store MESH_REGISTRY names, and it names none")
}
body, _, err = (artifacts.Store{Address: address}).GetTagged(ctx, snapshot.Repository, snapshot.Tag)
} else {
body, err = os.ReadFile(from)
}
if err != nil {
return snapshot.Facts{}, err
}
return snapshot.Decode(body)
}
// judgeChange is the gate: both compositions, and what the change did to each machine.
func judgeChange(ctx context.Context, in mergeCheckInput) (mergeVerdict, error) {
v := mergeVerdict{Judge: version, Facts: in.facts.Taken, Modules: map[string]string{}}
say := func(format string, args ...any) {
if in.say != nil {
fmt.Fprintf(in.say, format+"\n", args...)
}
}
base, err := shelfOfFacts(in.facts)
if err != nil {
return v, err
}
change := base
sources := sourcesOfFacts(in.facts)
// **What the change moves is the planner's answer** (reachOfMerge, novox/hq ADR 0238), asked of the
// snapshot as if the change were merged: the definitions composed with the change are those of the
// modules a merge of it would rebuild, and of the modules it adds — not every definition in the tree,
// which may differ from what the mesh holds for reasons that are not this change's (a module pinned
// at an older commit, a main the mesh has not built yet).
var reach *mergeReach
if in.repository != "" && len(in.changed) > 0 {
r, err := reachOfChange(in.facts, in.repository, in.changed, in.tree)
if err != nil {
v.Notes = append(v.Notes, "what a merge of the change would rebuild could not be worked out, so every "+
"definition in its tree is judged: "+err.Error())
} else {
reach = &r
}
}
if in.tree != "" {
var failures []string
change, failures, err = shelfWithChange(base, sources, in.facts, in.repository, in.tree, v.Modules,
scopeOfReach(reach, in.repository))
if err != nil {
return v, err
}
v.Failures = append(v.Failures, failures...)
}
say("judging %d module(s) changed against %d machine(s), as the snapshot of %s says they run",
len(v.Modules), len(in.facts.Machines), in.facts.Taken.Format(time.RFC3339))
// What no single machine shows: the rules between manifests, identities on the mesh's own longest
// name, and the data each module keeps. Only what the change adds is the change's.
for _, rule := range []func(catalogue.Shelf) []string{
catalogue.CatalogueProblems,
func(s catalogue.Shelf) []string { return catalogue.IdentityProblems(s, in.facts.Longest()) },
catalogue.DataProblems,
} {
was := rule(catalogue.Shelf(base))
for _, p := range rule(catalogue.Shelf(change)) {
if !slices.Contains(was, p) {
v.Failures = append(v.Failures, p)
}
}
}
// Both compositions, each in a store of its own.
say("composing every machine as the mesh is")
before, _, notes, err := composeEveryMachine(ctx, in, base, sources, nil)
if err != nil {
return v, fmt.Errorf("the mesh as the snapshot says it is could not be raised: %w", err)
}
v.Notes = append(v.Notes, notes...)
say("composing every machine with the change")
var added []string
for module, how := range v.Modules {
if how == "added" {
added = append(added, module)
}
}
sort.Strings(added)
after, tried, notes, err := composeEveryMachine(ctx, in, change, sources, added)
if err != nil {
return v, fmt.Errorf("the mesh with the change could not be raised: %w", err)
}
for _, n := range notes {
if !slices.Contains(v.Notes, n) {
v.Failures = append(v.Failures, n)
}
}
for _, m := range in.facts.Machines {
gm := mergeMachine{Name: m.Name, Described: m.Described(), Live: m.Declaration.Composes,
Base: before[m.Name], Change: after[m.Name]}
gm.Added, gm.Removed = difference(gm.Base.Resources, gm.Change.Resources)
v.Machines = append(v.Machines, gm)
switch {
case gm.Base.Composes && !gm.Change.Composes:
v.Failures = append(v.Failures, fmt.Sprintf("%s: nothing could be sent to it with this change — %s",
gm.Described, firstOr(newOnly(gm.Change.Problems, gm.Base.Problems), gm.Change.Problems)))
case !gm.Base.Composes && m.Declaration.Composes:
// The mesh composes it and the gate could not raise it as it is: a fact the snapshot does
// not carry. Said, and judged by what the change adds.
v.Notes = append(v.Notes, fmt.Sprintf("%s composes on the mesh and not as the snapshot raised it: %s — "+
"judged by what the change adds", gm.Described, firstOr(gm.Base.Problems, nil)))
if added := newOnly(gm.Change.Problems, gm.Base.Problems); len(added) > 0 {
v.Failures = append(v.Failures, fmt.Sprintf("%s: the change adds — %s", gm.Described, added[0]))
}
case !gm.Base.Composes:
v.Notes = append(v.Notes, fmt.Sprintf("%s does not compose on the mesh today either: %s",
gm.Described, firstOr(m.Declaration.Problems, gm.Base.Problems)))
}
for _, w := range newOnly(gm.Change.Withheld, gm.Base.Withheld) {
v.Failures = append(v.Failures, fmt.Sprintf("%s: the change leaves a consumer out of its grants — %s",
gm.Described, w))
}
for _, u := range newOnly(gm.Change.Unbound, gm.Base.Unbound) {
v.Failures = append(v.Failures, fmt.Sprintf("%s: the change leaves a credential bound elsewhere — %s",
gm.Described, u))
}
for module, why := range gm.Change.LeftOut {
if _, was := gm.Base.LeftOut[module]; !was {
v.Failures = append(v.Failures, fmt.Sprintf("%s: the change leaves %s out of its declaration — %s",
gm.Described, module, why))
}
}
}
// A module nobody runs yet, tried on a machine that could run it.
for _, module := range added {
t := tried[module]
for _, r := range t.Refused {
v.Failures = append(v.Failures, fmt.Sprintf("%s, assigned to %s, would be refused by its node-engine — %s",
module, t.Machine, r))
}
if t.Said != "" {
v.Notes = append(v.Notes, fmt.Sprintf("%s could not be tried on a machine: %s", module, t.Said))
}
}
if reach != nil {
w := widthOf(*reach)
v.Width = &w
if len(w.Unread) > 0 {
v.Notes = append(v.Notes, fmt.Sprintf("%s read by no module's build, and rebuild nothing (issue 280)",
readableList(w.Unread)))
}
if len(w.Modules) > wideRebuild {
v.Warnings = append(v.Warnings, fmt.Sprintf("a merge rebuilds %d module(s) in %d tier(s)",
len(w.Modules), w.Tiers))
}
}
sort.Strings(v.Failures)
v.Failures = slices.Compact(v.Failures)
switch {
case len(v.Failures) > 0:
v.Verdict = "fail"
v.Summary = fmt.Sprintf("%d problem(s) the change brings; the first: %s", len(v.Failures), v.Failures[0])
case len(v.Warnings) > 0:
v.Verdict = "warning"
v.Summary = v.Warnings[0]
default:
v.Verdict = "pass"
composing := 0
for _, m := range v.Machines {
if m.Change.Composes {
composing++
}
}
v.Summary = fmt.Sprintf("every machine composes with the change as it did without (%d of %d compose)",
composing, len(v.Machines))
}
return v, nil
}
// Report is the verdict as a person reads it, and as a pull request is told it.
func (v mergeVerdict) Report() string {
var b strings.Builder
fmt.Fprintf(&b, "merge gate: %s — %s\n", strings.ToUpper(v.Verdict), v.Summary)
fmt.Fprintf(&b, "judged by controller %s against the facts of %s\n", v.Judge, v.Facts.Format(time.RFC3339))
list := func(title string, items []string) {
if len(items) == 0 {
return
}
fmt.Fprintf(&b, "\n%s:\n", title)
for _, i := range items {
fmt.Fprintf(&b, " - %s\n", i)
}
}
list("fails", v.Failures)
list("warns", v.Warnings)
if len(v.Modules) > 0 {
var changed []string
for m, how := range v.Modules {
changed = append(changed, m+" ("+how+")")
}
sort.Strings(changed)
list("modules the change moves", changed)
}
var machines []string
for _, m := range v.Machines {
line := m.Described + ": "
switch {
case m.Change.Composes:
line += "composes"
default:
line += "does not compose"
}
if len(m.Added) > 0 || len(m.Removed) > 0 {
line += fmt.Sprintf("; the change adds %d resource(s) and removes %d", len(m.Added), len(m.Removed))
if len(m.Added)+len(m.Removed) <= 6 {
line += " (" + strings.Join(append(prefixed("+", m.Added), prefixed("-", m.Removed)...), ", ") + ")"
}
}
machines = append(machines, line)
}
list("machines", machines)
if v.Width != nil {
fmt.Fprintf(&b, "\na merge would rebuild %d module(s) in %d tier(s)", len(v.Width.Modules), v.Width.Tiers)
if len(v.Width.Modules) > 0 && len(v.Width.Modules) <= wideRebuild {
fmt.Fprintf(&b, ": %s", strings.Join(v.Width.Modules, ", "))
}
fmt.Fprintln(&b)
}
list("notes", v.Notes)
return b.String()
}
func prefixed(p string, items []string) []string {
out := make([]string, len(items))
for i, s := range items {
out[i] = p + s
}
return out
}
// shelfOfFacts is every module of the snapshot, as the mesh holds it — but the ones that came with the
// controller: those are this controller's own, provided when its store is raised.
func shelfOfFacts(f snapshot.Facts) (map[string]catalogue.Manifest, error) {
out := map[string]catalogue.Manifest{}
for _, m := range f.Modules {
if m.Provided {
continue
}
var manifest catalogue.Manifest
if err := json.Unmarshal(m.Manifest, &manifest); err != nil {
return nil, fmt.Errorf("the snapshot's %s is not a manifest: %w", m.Name, err)
}
out[m.Name] = manifest
}
return out, nil
}
// sourcesOfFacts is where each module of the snapshot is built from.
func sourcesOfFacts(f snapshot.Facts) map[string]inventory.Source {
out := map[string]inventory.Source{}
for _, m := range f.Modules {
if !m.Provided {
out[m.Name] = inventory.Source{Repository: m.Repository, Path: m.Path, BuiltFrom: m.Commit}
}
}
return out
}
// ignoredInTree are directories a module.json in is not a module of the repository.
var ignoredInTree = map[string]bool{".git": true, "vendor": true, "node_modules": true, "testdata": true,
"examples": true, "fixtures": true}
// shelfWithChange is the base shelf with every module of the change's tree in place of the mesh's: each
// read by this controller's strict parser — the one the mesh runs, for a change to a catalogue — resolved
// with stand-in builds, and read again as registration reads a built one. A module the repository held
// and the tree no longer has is removed. Answers the shelf and what the tree itself fails.
//
// `scope` is the directories of the modules a merge of the change would rebuild or add, as the planner says
// (scopeOfReach); a definition elsewhere in the tree is left as the mesh holds it. Nil judges every one.
func shelfWithChange(base map[string]catalogue.Manifest, sources map[string]inventory.Source, f snapshot.Facts,
repository, tree string, moved map[string]string, scope map[string]bool) (map[string]catalogue.Manifest, []string, error) {
inScope := func(dir string) bool { return scope == nil || scope[dir] }
out := make(map[string]catalogue.Manifest, len(base))
for k, m := range base {
out[k] = m
}
var failures []string
found := map[string]string{} // module → its directory in the tree
err := filepath.WalkDir(tree, func(path string, d fs.DirEntry, err error) error {
if err != nil {
return err
}
if d.IsDir() {
if path != tree && (ignoredInTree[d.Name()] || strings.HasPrefix(d.Name(), ".")) {
return filepath.SkipDir
}
return nil
}
if d.Name() != "module.json" {
return nil
}
dir, err := filepath.Rel(tree, filepath.Dir(path))
if err != nil {
return err
}
if dir == "." {
dir = ""
}
raw, err := os.ReadFile(path)
if err != nil {
return err
}
m, err := catalogue.ParseManifest(raw)
if err != nil {
if inScope(dir) {
failures = append(failures, fmt.Sprintf("%s: the controller judging this (%s) refuses it — %s",
orRoot(dir), version, oneLine(err.Error())))
}
return nil
}
if was, twice := found[m.Module]; twice {
if inScope(dir) || inScope(was) {
failures = append(failures, fmt.Sprintf("%s and %s both define %s: two definitions of one module",
orRoot(was), orRoot(dir), m.Module))
}
return nil
}
found[m.Module] = dir
if !inScope(dir) {
return nil
}
if s, held := sources[m.Module]; held && s.Repository != "" && !sameRepository(s.Repository,
link.SourceMoved{Owner: ownerOf(repository), Repo: repoOf(repository)}) {
failures = append(failures, fmt.Sprintf("%s defines %s, which the mesh builds from %s: two definitions "+
"of one module name (issue 239)", orRoot(dir), m.Module, s.Repository))
return nil
}
resolved, err := standInBuild(m)
if err == nil {
// As registration reads a built manifest: strictly, again (takeIn).
var raw []byte
if raw, err = json.Marshal(resolved); err == nil {
resolved, err = catalogue.ParseManifest(raw)
}
}
if err == nil {
err = namesNoInstallation(resolved)
}
if err != nil {
failures = append(failures, fmt.Sprintf("%s: %s is not registrable — %s", orRoot(dir), m.Module,
oneLine(err.Error())))
return nil
}
if was, held := base[m.Module]; !held {
moved[m.Module] = "added"
} else if builtAlike(was) == builtAlike(resolved) {
// Unchanged but for what a build pins: the mesh's own build stands, so the machines that
// run it compose exactly as they do.
return nil
} else {
moved[m.Module] = "changed"
}
out[m.Module] = resolved
return nil
})
if err != nil {
return nil, nil, err
}
// What the repository held and the tree no longer defines is removed — refused while a machine runs it.
running := map[string][]string{}
for _, mc := range f.Machines {
for _, a := range mc.Assigned {
running[a] = append(running[a], mc.Described())
}
}
for name, s := range sources {
if s.Repository == "" || !sameRepository(s.Repository, link.SourceMoved{Owner: ownerOf(repository), Repo: repoOf(repository)}) {
continue
}
if _, still := found[name]; still {
continue
}
if !inScope(strings.Trim(s.Path, "/")) {
continue
}
delete(out, name)
moved[name] = "removed"
if on := running[name]; len(on) > 0 {
failures = append(failures, fmt.Sprintf("%s is removed from %s, and %s run(s) it: unassign it first (ADR 0236)",
name, repository, readableList(on)))
}
}
return out, failures, nil
}
func orRoot(dir string) string {
if dir == "" {
return "the repository's root"
}
return dir
}
func ownerOf(repository string) string {
owner, _, _ := strings.Cut(repository, "/")
return owner
}
func repoOf(repository string) string {
_, repo, _ := strings.Cut(repository, "/")
return repo
}
// builtAlike is a resolved manifest with what a build pins taken out — the digests and the references
// they are reached by, and the version a path is named for — so a definition the change leaves alone reads
// the same as the build of it the mesh holds.
func builtAlike(m catalogue.Manifest) string {
raw, err := json.Marshal(m)
if err != nil {
return ""
}
var doc any
if err := json.Unmarshal(raw, &doc); err != nil {
return ""
}
var walk func(any) any
walk = func(v any) any {
switch t := v.(type) {
case map[string]any:
for k, e := range t {
switch strings.ToLower(k) {
case "image", "source", "digest", "launchers":
t[k] = "pinned"
default:
t[k] = walk(e)
}
}
return t
case []any:
for i, e := range t {
t[i] = walk(e)
}
return t
case string:
if before, _, found := strings.Cut(t, "/versions/"); found {
return before + "/versions/pinned"
}
return t
}
return v
}
out, err := json.Marshal(walk(doc))
if err != nil {
return ""
}
return string(out)
}
// standInBuild resolves a manifest as a build would, with artifacts nobody built: a digest made from the
// module's and the artifact's names, in the store's own vocabulary. Composition and validation read the
// shape of a reference, never its bytes.
func standInBuild(m catalogue.Manifest) (catalogue.Manifest, error) {
if m.Build == nil {
return m.Resolve(nil)
}
var built []catalogue.Built
for _, a := range m.Build.Artifacts {
sum := sha256.Sum256([]byte("gate\x00" + m.Module + "\x00" + a.Name))
digest := fmt.Sprintf("sha256:%x", sum)
b := catalogue.Built{Name: a.Name, Kind: a.Kind, Digest: digest}
switch a.Kind {
case catalogue.ArtifactImage, catalogue.ArtifactUpstream:
b.Reference = catalogue.ArtifactStoreScheme + m.Module + "/" + a.Name + "@" + digest
default:
b.Reference = catalogue.ArtifactStoreScheme + m.Module + "/" + a.Name + "/blobs/" + digest
}
built = append(built, b)
}
return m.Resolve(built)
}
// composeEveryMachine raises a throwaway store from the snapshot with this shelf, composes every machine
// twice — the first time as a push does, making each credential, so the second has every consumer's
// grant to compose — and answers each machine's composition. Notes are what of the snapshot could not
// be raised: a refusal to register a module, keep a setting or a pin.
func composeEveryMachine(ctx context.Context, in mergeCheckInput, shelf map[string]catalogue.Manifest,
sources map[string]inventory.Source, trials []string) (map[string]mergeComposed, map[string]mergeTrial, []string, error) {
open, drop, err := throwawayStores(ctx, in.admin)
if err != nil {
return nil, nil, nil, err
}
defer drop()
notes, err := raiseFromFacts(ctx, open, in.facts, shelf, sources)
if err != nil {
return nil, nil, notes, err
}
gens, gensErr := generators(ctx, open)
out := map[string]mergeComposed{}
for pass := 0; pass < 2; pass++ {
for _, m := range in.facts.Machines {
if gensErr != nil {
out[m.Name] = mergeComposed{Problems: []string{"the private network cannot be computed: " +
oneLine(gensErr.Error())}}
continue
}
declared, problems, err := composedAndValidated(ctx, open, m.Name, gens, Allocating)
if pass == 0 {
continue
}
c := mergeComposed{}
if err != nil {
c.Problems = []string{oneLine(err.Error())}
out[m.Name] = c
continue
}
c.Composes = len(problems) == 0
c.Problems = problems
c.Resources = resourceNames(declared.Resources)
c.LeftOut = declared.leftOutWhy
for _, o := range declared.withheld {
c.Withheld = append(c.Withheld, o.String())
}
for _, u := range declared.unbound {
c.Unbound = append(c.Unbound, u.String())
}
sort.Strings(c.Withheld)
sort.Strings(c.Unbound)
out[m.Name] = c
}
if ctx.Err() != nil {
return nil, nil, notes, ctx.Err()
}
}
tried := map[string]mergeTrial{}
if gensErr == nil {
for _, module := range trials {
tried[module] = tryAssigning(ctx, open, in.facts, shelf[module], gens)
}
}
return out, tried, notes, nil
}
// mergeTrial is a module no machine runs yet, assigned for a moment to one that could run it.
type mergeTrial struct {
Machine string `json:"machine,omitempty"`
Refused []string `json:"refused,omitempty"`
Composes bool `json:"composes"`
Said string `json:"said,omitempty"`
}
// tryAssigning composes a module nobody runs yet on the first machine that has what it needs, and says
// what the node-engine would refuse of it there (issue 236: the refusal came only from the first live
// machine it was assigned to). The assignment is taken back; nothing else in the store moves.
func tryAssigning(ctx context.Context, open *stores, f snapshot.Facts, m catalogue.Manifest,
gens map[string]catalogue.Generator) mergeTrial {
for _, machine := range f.Machines {
has := map[string]bool{}
for _, c := range machine.Capabilities {
has[c.Name] = c.Present
}
fits := true
for _, c := range m.Capabilities {
fits = fits && has[c]
}
if !fits {
continue
}
if _, err := open.inventory.Assign(ctx, machine.Name, m.Module); err != nil {
return mergeTrial{Machine: machine.Described(), Said: oneLine(err.Error())}
}
_, problems, err := composedAndValidated(ctx, open, machine.Name, gens, Allocating)
_ = open.inventory.Unassign(ctx, machine.Name, m.Module)
t := mergeTrial{Machine: machine.Described()}
if err != nil {
t.Said = oneLine(err.Error())
return t
}
for _, p := range problems {
if strings.Contains(p, `"`+m.Module+".") || strings.Contains(p, m.Module+":") {
t.Refused = append(t.Refused, p)
}
}
t.Composes = len(problems) == 0
return t
}
return mergeTrial{Said: "no machine of the mesh has what it needs: " + strings.Join(m.Capabilities, ", ")}
}
// throwawayStores creates fresh databases for every context the controller keeps in the PostgreSQL at
// admin, brings each to this controller's schema, and answers them opened, with what drops them again.
func throwawayStores(ctx context.Context, admin string) (*stores, func(), error) {
suffix := fmt.Sprintf("%d", time.Now().UnixNano()%1_000_000_000)
conn, err := pgx.Connect(ctx, admin)
if err != nil {
return nil, nil, fmt.Errorf("cannot reach the throwaway PostgreSQL: %w", err)
}
defer conn.Close(ctx)
cut := strings.LastIndex(admin, "/")
if cut < 0 {
return nil, nil, fmt.Errorf("%q is not a PostgreSQL URL", admin)
}
query := ""
if q := strings.Index(admin[cut:], "?"); q >= 0 {
query = admin[cut+q:]
}
var names []string
// The stores are found by the environment, so what it said before is what it says after: a gate run
// inside a process with stores of its own — a test, a controller — leaves them as they were.
restore := map[string]*string{}
for _, c := range held {
if was, set := os.LookupEnv(store.Variable(c.name)); set {
restore[c.name] = &was
} else {
restore[c.name] = nil
}
}
drop := func() {
for name, was := range restore {
if was == nil {
_ = os.Unsetenv(store.Variable(name))
} else {
_ = os.Setenv(store.Variable(name), *was)
}
}
c, err := pgx.Connect(context.Background(), admin)
if err != nil {
return
}
defer c.Close(context.Background())
for _, n := range names {
_, _ = c.Exec(context.Background(), "drop database if exists "+n+" with (force)")
}
}
for _, c := range held {
name := "gate_" + c.name + "_" + suffix
if _, err := conn.Exec(ctx, "create database "+name); err != nil {
drop()
return nil, nil, fmt.Errorf("cannot create %s: %w", name, err)
}
names = append(names, name)
url := admin[:cut] + "/" + name + query
if err := os.Setenv(store.Variable(c.name), url); err != nil {
drop()
return nil, nil, err
}
migrations, err := c.migrations()
if err != nil {
drop()
return nil, nil, err
}
s, err := store.Open(ctx, c.name)
if err != nil {
drop()
return nil, nil, err
}
_, err = s.Migrate(ctx, migrations)
s.Close()
if err != nil {
drop()
return nil, nil, fmt.Errorf("%s does not migrate: %w", c.name, err)
}
}
open, err := openStores(ctx)
if err != nil {
drop()
return nil, nil, err
}
for _, m := range provided {
if err := open.inventory.Provide(ctx, m); err != nil {
open.Close()
drop()
return nil, nil, err
}
}
if _, err := open.inventory.SeedSeats(ctx, catalogue.DefaultSeats()); err != nil {
open.Close()
drop()
return nil, nil, err
}
return open, func() { open.Close(); drop() }, nil
}
// raiseFromFacts puts the snapshot's mesh into a fresh store through the controller's own records:
// every module of the shelf registered, every machine added where it is with what it reported, a key
// for each standing in for the one it holds, its assignments, the seats it holds, its pins and its
// settings, and a stand-in for every secret a person gave.
func raiseFromFacts(ctx context.Context, open *stores, f snapshot.Facts, shelf map[string]catalogue.Manifest,
sources map[string]inventory.Source) ([]string, error) {
inv := open.inventory
var notes []string
names := make([]string, 0, len(shelf))
for name := range shelf {
names = append(names, name)
}
sort.Strings(names)
for _, name := range names {
s := sources[name]
if err := inv.RegisterModule(ctx, shelf[name], inventory.Source{Repository: s.Repository, Path: s.Path,
BuiltFrom: s.BuiltFrom}); err != nil {
notes = append(notes, fmt.Sprintf("%s is not registered: %s", name, oneLine(err.Error())))
}
}
for i, m := range f.Machines {
node, err := inv.AddNodeAs(ctx, m.Name, m.Adopted)
if err != nil {
return notes, fmt.Errorf("%s: %w", m.Name, err)
}
if m.OnNetwork {
endpoint := ""
if m.Public {
endpoint = fmt.Sprintf("192.0.2.%d:51820", i+1)
}
if err := inv.SetPlace(ctx, m.Name, endpoint, m.Site, m.Hub, fmt.Sprintf("10.99.%d.%d", i/250, i%250+1)); err != nil {
return notes, fmt.Errorf("%s: %w", m.Name, err)
}
}
var caps []map[string]any
for _, c := range m.Capabilities {
caps = append(caps, map[string]any{"name": c.Name, "present": c.Present, "detail": c.Detail})
}
if err := inv.RecordProfile(ctx, node.ID, map[string]any{"architecture": m.Architecture,
"kernel": m.Kernel, "capabilities": caps}); err != nil {
return notes, err
}
for _, record := range []func(context.Context, string, string) error{inv.RecordSealingKey, inv.RecordOverlayKey} {
key, err := standInKey()
if err != nil {
return notes, err
}
if err := record(ctx, node.ID, key); err != nil {
return notes, err
}
}
if m.Account != "" {
if err := inv.SetAccount(ctx, m.Name, m.Account, m.AccountHome); err != nil {
return notes, err
}
}
if m.PublicDomain != "" {
if err := inv.SetPublicDomain(ctx, m.Name, m.PublicDomain); err != nil {
return notes, err
}
}
}
for _, m := range f.Machines {
for _, module := range m.Assigned {
if _, err := inv.Assign(ctx, m.Name, module); err != nil {
notes = append(notes, fmt.Sprintf("%s cannot be assigned to %s: %s", module, m.Described(),
oneLine(err.Error())))
}
}
}
for _, s := range f.Seats {
for _, h := range s.Holders {
var err error
if s.Scope == catalogue.ScopeMesh && len(s.Holders) > 1 {
err = inv.AddSeatHolder(ctx, s.Name, s.Scope, h.Machine, h.Module)
} else {
err = inv.HoldSeat(ctx, s.Name, s.Scope, h.Machine, h.Module)
}
if err != nil {
notes = append(notes, fmt.Sprintf("%s's hold of %s is not kept: %s", h.Machine, s.Name, oneLine(err.Error())))
}
}
}
for _, s := range f.Settings {
if err := inv.SetSettings(ctx, "", s.Module, s.Values); err != nil {
notes = append(notes, fmt.Sprintf("the mesh's settings of %s are not kept: %s", s.Module, oneLine(err.Error())))
}
}
for _, m := range f.Machines {
for _, p := range m.Pins {
if err := inv.PinProvision(ctx, m.Name, p.Provision, p.Machine, p.Module); err != nil {
notes = append(notes, fmt.Sprintf("%s's pin of %s is not kept: %s", m.Described(), p.Provision,
oneLine(err.Error())))
}
}
for _, s := range m.Settings {
if err := inv.SetSettings(ctx, m.Name, s.Module, s.Values); err != nil {
notes = append(notes, fmt.Sprintf("%s's settings of %s are not kept: %s", m.Described(), s.Module,
oneLine(err.Error())))
}
}
for _, a := range m.Accepted {
standIn := fmt.Sprintf("gate-stand-in-%x", sha256.Sum256([]byte(m.Name+a.Module+a.Name+a.Provider+a.Local)))[:40]
var err error
if a.Provider == "" {
err = inv.AcceptSecretForModule(ctx, m.Name, a.Module, a.Name, standIn)
} else {
err = inv.AcceptSecretForPair(ctx, a.Name, m.Name, a.Module, a.Provider, a.Local, standIn)
}
if err != nil {
notes = append(notes, fmt.Sprintf("%s's given %s for %s is not kept: %s", m.Described(), a.Name,
a.Module, oneLine(err.Error())))
}
}
}
return notes, nil
}
// standInKey is a key a machine could have reported; its private half is never kept.
func standInKey() (string, error) {
k, err := ecdh.X25519().GenerateKey(rand.Reader)
if err != nil {
return "", err
}
return base64.StdEncoding.EncodeToString(k.PublicKey().Bytes()), nil
}
// reachOfChange is what a merge of the change would move and build — **the planner's own answer**
// (reachOfMerge), asked of the modules, sources, reads and edges the snapshot carries, never a mapping of
// the gate's own: what a merge touches and what follows it are decided in one place, so the gate and the
// plan a merge makes cannot disagree (novox/hq ADR 0238).
func reachOfChange(f snapshot.Facts, repository string, paths []string, tree string) (mergeReach, error) {
owner, repo, found := strings.Cut(repository, "/")
if !found {
return mergeReach{}, fmt.Errorf("%q is not owner/repository", repository)
}
m := link.SourceMoved{Owner: owner, Repo: repo, Base: "main", Commit: "gate", Paths: paths}
// Which changed directories hold a module, read from the change's tree as the forge's announcer reads
// them at the head (issue 278): what the planner is told of a module the mesh does not hold yet.
if tree != "" {
m.ModuleDirs, m.ModuleDirsSaid = moduleDirsIn(tree, paths), true
for _, p := range paths {
if _, err := os.Stat(filepath.Join(tree, p)); os.IsNotExist(err) {
m.Removed = append(m.Removed, p)
}
}
}
var entries []inventory.Entry
read := map[string][]inventory.ReadRepository{}
for _, mod := range f.Modules {
var manifest catalogue.Manifest
if err := json.Unmarshal(mod.Manifest, &manifest); err != nil {
return mergeReach{}, err
}
entries = append(entries, inventory.Entry{Manifest: manifest, Provided: mod.Provided,
Source: inventory.Source{Repository: mod.Repository, Path: mod.Path, BuiltFrom: mod.Commit}})
for _, r := range mod.Reads {
read[mod.Name] = append(read[mod.Name], inventory.ReadRepository{Repository: r})
}
}
var edges []inventory.Edge
for _, e := range f.Edges {
edges = append(edges, inventory.Edge{From: e.From, To: e.To, Kind: e.Kind})
}
return reachOfMerge(m, entries, read, edges), nil
}
// widthOf is how wide the rebuild of a reach is.
func widthOf(r mergeReach) mergeWidthOf {
w := mergeWidthOf{Tiers: len(r.Plan.Tiers), Unread: r.Unread}
for name := range r.Plan.Modules {
w.Modules = append(w.Modules, name)
}
sort.Strings(w.Modules)
return w
}
// scopeOfReach is the directories, in the change's repository, of the modules a reach rebuilds or adds;
// nil — every definition judged — when there is no reach.
func scopeOfReach(r *mergeReach, repository string) map[string]bool {
if r == nil {
return nil
}
scope := map[string]bool{}
same := link.SourceMoved{Owner: ownerOf(repository), Repo: repoOf(repository)}
for _, e := range append(append([]inventory.Entry{}, r.Touched...), r.Deleted...) {
if sameRepository(e.Source.Repository, same) {
scope[strings.Trim(e.Source.Path, "/")] = true
}
}
for _, d := range r.Added {
if d == "." {
d = ""
}
scope[d] = true
}
return scope
}
// moduleDirsIn are the directories above the changed paths, never the root, that hold a module.json in
// the tree.
func moduleDirsIn(tree string, paths []string) []string {
seen := map[string]bool{}
var out []string
for _, p := range paths {
for dir := filepath.Dir(filepath.Clean(p)); dir != "." && dir != "/" && dir != ""; dir = filepath.Dir(dir) {
if seen[dir] {
continue
}
seen[dir] = true
if _, err := os.Stat(filepath.Join(tree, dir, "module.json")); err == nil {
out = append(out, filepath.ToSlash(dir))
}
}
}
sort.Strings(out)
return out
}
// difference is what b has that a has not, and what a has that b has not.
func difference(a, b []string) (added, removed []string) {
for _, x := range b {
if !slices.Contains(a, x) {
added = append(added, x)
}
}
for _, x := range a {
if !slices.Contains(b, x) {
removed = append(removed, x)
}
}
return added, removed
}
// newOnly is what now has that was had not.
func newOnly(now, was []string) []string {
var out []string
for _, x := range now {
if !slices.Contains(was, x) {
out = append(out, x)
}
}
return out
}
func firstOr(items, otherwise []string) string {
if len(items) > 0 {
return items[0]
}
if len(otherwise) > 0 {
return otherwise[0]
}
return "it says nothing more"
}