Merge pull request 'Say why assign finds no module, and keep an assignment pending on its build (hq issue 325)' (#150) from fix/assign-says-why-a-module-is-not-there into main

This commit was merged in pull request #150.
This commit is contained in:
2026-10-08 14:58:05 +00:00
21 changed files with 2425 additions and 16 deletions
+41
View File
@@ -7,8 +7,10 @@ import (
"slices"
"sort"
"strings"
"time"
"github.com/novox/mesh-controller/internal/catalogue"
"github.com/novox/mesh-controller/internal/inventory"
)
// The things the mesh can be asked to do, separated from how it was asked.
@@ -41,6 +43,11 @@ import (
// It costs a resolution per machine. Assignment is a person typing a command, and being told which
// machines this just blocked is worth more than the milliseconds.
func assign(ctx context.Context, open *stores, node string, modules ...string) (string, error) {
return assignWith(ctx, open, node, assignOptions{}, modules...)
}
// assignWith is assign asked with its options: build, for a module known and not built (novox/hq issue 325).
func assignWith(ctx context.Context, open *stores, node string, opts assignOptions, modules ...string) (string, error) {
if len(modules) == 0 {
return "", fmt.Errorf("assign %s names no module", node)
}
@@ -51,6 +58,29 @@ func assign(ctx context.Context, open *stores, node string, modules ...string) (
return "", err
}
defer release()
// **A module the catalogue does not hold is looked for further before it is refused** (novox/hq issue
// 325): a build in flight keeps the assignment pending, a module known and not built is said with how to
// build it, and only a name the mesh never heard of is refused as unknown. Several modules in one act
// are judged together (ADR 0207), so an act naming one not registered is not made in part.
if missing, err := notInCatalogue(ctx, open, modules); err != nil {
return "", err
} else if len(missing) > 0 {
if len(modules) == 1 {
return notRegistered(ctx, open, node, modules[0], opts)
}
var lines []string
for _, m := range missing {
where, err := whereIs(ctx, open.inventory, m, time.Now())
if err != nil {
return "", err
}
lines = append(lines, where.said)
}
return "", fmt.Errorf("%w: nothing was assigned, since modules assigned together are judged together "+
"(ADR 0207) and %s not registered:\n %s\n assign them together once each is registered, or each "+
"alone to keep it pending on its build", inventory.ErrNoSuchModule, strings.Join(missing, ", "),
strings.Join(lines, "\n "))
}
// **The one assignment refused for what the node lacks** (novox/hq ADR 0207). Everything else
// an assignment leaves unresolved is kept, because assignment is not an ordering; a module whose
// resources are applied through a seat nothing on the node holds is refused, because that order
@@ -200,6 +230,17 @@ func unassign(ctx context.Context, open *stores, node string, modules ...string)
for _, a := range assigned {
runs[a] = true
}
// A module only pending on the node (novox/hq issue 325) is withdrawn, when it is the act's one module:
// that is the undo of an assignment kept while its build runs.
if len(modules) == 1 && !runs[modules[0]] {
said, withdrawn, err := withdrawPending(ctx, open.inventory, node, modules[0])
if err != nil {
return "", err
}
if withdrawn {
return said, nil
}
}
for _, module := range modules {
if !runs[module] {
return "", fmt.Errorf("%s is not assigned to %s", module, node)
+62 -2
View File
@@ -399,15 +399,49 @@ func buildBehind(ctx context.Context, wait time.Duration) error {
//
// Separated from the command so `--behind` can walk a list without a second path to the same act.
func buildOne(ctx context.Context, source buildSource, path, ref string, wait time.Duration) error {
_, err := buildOneAsked(ctx, source, path, ref, wait, false)
_, err := buildOneAsked(ctx, source, path, ref, wait, false, "build")
return err
}
// recordAsked keeps a build request once it is asked (novox/hq issue 325), with who asked it, in a store
// opened for the purpose: the asks reach here from processes that hold none.
func recordAsked(ctx context.Context, r inventory.BuildRequest) {
open, err := openStores(ctx)
if err != nil {
fmt.Fprintf(os.Stderr, "build %s is asked and not kept as a build request: %v\n", r.ID, err)
return
}
defer open.Close()
recordBuildRequest(ctx, open.inventory, r)
}
// markWaitFailed says what became of a kept build request whose waited ask failed: never handed over, so
// nothing runs; or handed over and no longer waited for, so asked with its outcome unknown — still read as in
// flight until its outcome or its bound (novox/hq issue 325).
func markWaitFailed(ctx context.Context, id string, why error) {
open, err := openStores(ctx)
if err != nil {
fmt.Fprintf(os.Stderr, "build %s could not be marked: %v\n", id, err)
return
}
defer open.Close()
mark := open.inventory.MarkOutcomeUnknown
if errors.Is(why, link.ErrNotHandedOver) {
mark = open.inventory.MarkNotAsked
}
if err := mark(ctx, id, why.Error()); err != nil {
fmt.Fprintf(os.Stderr, "build %s could not be marked: %v\n", id, err)
}
}
// buildOneAsked is buildOne answering the id it asked with — what a plan keeps to match the outcome
// by (novox/hq ADR 0219) — and, for an ask not waited for, optionally a dry run: built and looked
// at, never taken in (issue 240), which is what `replay` asks unless told to register.
//
// asker is who asked, for the build request kept once the ask is made (novox/hq issue 325); empty for an
// asker that keeps the request itself, with its own name, once it knows the ask was made.
func buildOneAsked(ctx context.Context, source buildSource, path, ref string, wait time.Duration,
dryRun bool) (string, error) {
dryRun bool, asker string) (string, error) {
if dryRun && wait != 0 {
return "", errors.New("a dry run waited for is `build --dry-run`")
}
@@ -462,6 +496,12 @@ func buildOneAsked(ctx context.Context, source buildSource, path, ref string, wa
}
defer ask.Close()
fmt.Printf(" of %s\n", seat)
// Kept once it is asked, so `assign` knows a module is coming while its build runs (novox/hq issue 325);
// never before, so an ask that failed never reads as a build in flight. A dry run registers nothing, and
// is not kept.
keep := !dryRun && asker != ""
asked := inventory.BuildRequest{ID: request.ID, Repository: source.Repository, Seat: source.Seat, Path: path,
Ref: ref, For: asker}
if wait == 0 {
// Asked and not waited for (novox/hq issue 176): the outcome is the role's event, and the
@@ -471,6 +511,9 @@ func buildOneAsked(ctx context.Context, source buildSource, path, ref string, wa
if err := ask.Ask(ctx, request); err != nil {
return "", err
}
if keep {
recordAsked(ctx, asked)
}
if dryRun {
fmt.Printf("asked as a dry run, not waited for: `builds --log %s` follows it as it runs; "+
"its outcome is not taken in\n", request.ID)
@@ -481,8 +524,16 @@ func buildOneAsked(ctx context.Context, source buildSource, path, ref string, wa
return request.ID, nil
}
// Waited for: kept while it is waited for, since that can take minutes, and marked when the hand-over
// failed (not asked) or the wait did (asked, outcome unknown) — an outcome heard later is the last word.
if keep {
recordAsked(ctx, asked)
}
result, err := ask.Submit(ctx, request, wait)
if err != nil {
if keep {
markWaitFailed(ctx, request.ID, err)
}
return request.ID, err
}
@@ -492,6 +543,10 @@ func buildOneAsked(ctx context.Context, source buildSource, path, ref string, wa
}
defer open.Close()
manifest, kept, err := takeIn(ctx, open.inventory, result)
// The assignments pending on this build, made or ended (novox/hq issue 325).
for _, line := range settlePendingOnBuild(ctx, open, result, manifest.Module, err) {
fmt.Printf(" %s\n", line)
}
if err != nil {
return request.ID, err
}
@@ -771,6 +826,11 @@ type answers struct {
// handed to the operator; healsUnread why it could not be read.
heals *healsCount
healsUnread string
// pending is every assignment waiting for its module's build, and every one ended in the last day
// (novox/hq issue 325); pendingUnread why they could not be read. One waiting, or one ended other than
// applied, is not well: an assignment somebody made is not made yet, or will not be.
pending []inventory.PendingAssignment
pendingUnread string
}
// heldBy is every artifact this mesh has built, for a build that may need one as its base.
+9
View File
@@ -371,6 +371,15 @@ func (b builds) Built(ctx context.Context, result link.BuildResult) error {
return nil
}
manifest, _, err := takeIn(ctx, b.inv, result)
// The assignments pending on this build, made or ended (novox/hq issue 325), said in the daemon's log
// after what became of the build.
if b.open != nil {
defer func() {
for _, line := range settlePendingOnBuild(ctx, b.open, result, manifest.Module, err) {
fmt.Printf("%s: %s\n", result.ID, line)
}
}()
}
// When it was asked, so a plan takes as its outcome only a build asked for it or after it
// (novox/hq 04-ISSUES/219). Zero when the id does not say.
asked, _ := link.BuildAskedAt(result.ID)
+14 -1
View File
@@ -354,9 +354,20 @@ func moduleCommand(ctx context.Context, args []string) error {
func assignCommand(ctx context.Context, verb string, args []string) error {
// Several modules in one act (novox/hq ADR 0207): holders that depend on each other — the
// service manager and the package manager — can only go on, or come off, together.
set := flag.NewFlagSet(verb, flag.ContinueOnError)
// A module known and not built: ask for its build, and keep the assignment pending on it (novox/hq
// issue 325). Only assign reads it.
build := set.Bool("build", false, "for a module known and not built: ask for its build, and assign it when it registers")
args, err := parseAround(set, args)
if err != nil {
return err
}
if len(args) < 2 {
return fmt.Errorf("%s <node> <module> [<module>…]", verb)
}
if *build && verb == "unassign" {
return errors.New("unassign takes no --build")
}
open, err := openStores(ctx)
if err != nil {
return err
@@ -365,7 +376,9 @@ func assignCommand(ctx context.Context, verb string, args []string) error {
// The act itself is in acts.go, so the command API refuses exactly what this refuses
// (novox/hq ADR 0035). What differs between the surfaces is how the answer is printed.
act := assign
act := func(ctx context.Context, open *stores, node string, modules ...string) (string, error) {
return assignWith(ctx, open, node, assignOptions{Build: *build}, modules...)
}
if verb == "unassign" {
act = unassign
}
+749
View File
@@ -0,0 +1,749 @@
package main
import (
"context"
"errors"
"fmt"
"os"
"sort"
"strings"
"time"
"github.com/novox/mesh-controller/internal/conditions"
"github.com/novox/mesh-controller/internal/inventory"
"github.com/novox/mesh-controller/internal/link"
)
// A module not registered yet, and the assignment that waits for it (novox/hq issue 325, ADR 0261).
//
// On 2026-10-08 the catalogue's pull request adding `sensors` merged; a minute later `assign g14 sensors`
// answered "no module of that name: sensors", and a few minutes later the same call worked. The merge had
// asked for the build (issue 300) and the build had not finished. The answer read as "you forgot to register
// it", and nothing in it said a build was on its way.
//
// So an assignment of a module the catalogue does not hold looks further before it refuses, and says which
// of three cases it is in:
//
// - **a build is in flight**: a build request for the module's directory has no outcome yet, or its build
// succeeded a moment ago and is being registered. The assignment is kept as a **pending assignment**,
// which the controller makes when the build registers the module (ADR 0261).
// - **known, not built**: the controller asked for the directory before, and the build failed, was not
// registered, said nothing within its bound, or could not be asked. Said, with `build` and assign's
// own build argument, which asks for the build and keeps the assignment pending on it.
// - **unknown**: no module registered and no build request kept by that name. Refused as before, with the
// closest names.
//
// **Pending by default while a build is in flight** (ADR 0261), without an argument to ask for it: an
// assignment is what a person meant and is kept even when its machine does not resolve (acts.go), and
// `unassign` withdraws it. Asking for a second act once the build lands is the two acts by hand issue 300
// removed. A build is never asked for by default: it runs on another machine, and a directory whose last
// build failed is a fault to look at before it is a thing to retry.
//
// **Settling is the controller's tick, never a read** (settlingPending): `status`, the board and the summary
// only read the rows.
// buildRequestBound is how long a build request may go without an outcome before it is no longer read as in
// flight, and a pending assignment waiting for it expires. Generous: a build waits behind others on the
// build seat, and a pending assignment that waits a little longer costs nothing.
const buildRequestBound = 2 * time.Hour
// registerGrace is how long a build heard as built is read as still in flight while it is not registered:
// the outcome is recorded first and registered after it, in the same take-in. Past it, the build was
// recorded and not registered, and says so.
const registerGrace = 5 * time.Minute
// pendingShownFor is how long an ended pending assignment is listed in `status`.
const pendingShownFor = 24 * time.Hour
// settleEvery is how often the controller settles the pending assignments.
const settleEvery = time.Minute
// claimStaleAfter is how long a pending assignment may be held as applying before the tick reads the claim
// as left by a controller that stopped, and settles it from what the mesh holds.
const claimStaleAfter = 5 * time.Minute
// kindPendingEnded is a pending assignment that ended without being made: expired or refused.
const kindPendingEnded = "assignment-not-made"
// assignOptions are what an assignment was asked with beside its machine and modules.
type assignOptions struct {
// Build asks for the build of a module known and not built, and keeps the assignment pending on it.
Build bool
}
// recordBuildRequest keeps a build request so `assign` can find it, once it is asked. Said, never fatal: the
// build was asked.
func recordBuildRequest(ctx context.Context, inv *inventory.Inventory, r inventory.BuildRequest) {
if err := inv.RecordBuildRequest(ctx, r); err != nil {
fmt.Fprintf(os.Stderr, "build %s was asked, and could not be kept as a build request, so `assign` "+
"will not know it is coming: %v\n", r.ID, err)
}
}
// The cases of a module the catalogue does not hold (novox/hq issue 325).
const (
unknownModule = iota
buildInFlight
knownNotBuilt
)
// notHeld is where a module the catalogue does not hold stands: its case, the build request it was read
// from, and the sentence that says it.
type notHeld struct {
kind int
request inventory.RequestOutcome
said string
}
// whereIs looks for a module the catalogue does not hold among the build requests the controller keeps.
func whereIs(ctx context.Context, inv *inventory.Inventory, module string, now time.Time) (notHeld, error) {
requests, err := inv.RequestsNamed(ctx, module)
if err != nil {
return notHeld{}, err
}
if r, ok := inFlight(requests, now); ok {
being := "being built"
if r.Heard {
being = "built and being registered"
}
return notHeld{kind: buildInFlight, request: r, said: fmt.Sprintf("%s is not registered yet: it is %s "+
"from %s since %s, as build %s; it can be assigned once that build registers it", module, being,
requestSource(r), clock(r.At), r.ID)}, nil
}
if len(requests) > 0 {
return notHeld{kind: knownNotBuilt, request: requests[0], said: fmt.Sprintf("%s is known and not "+
"registered: %s", module, whyNotBuilt(requests[0], now))}, nil
}
// Only what was looked at is claimed: the catalogue, and the build requests the controller keeps.
said := fmt.Sprintf("%v: %s — the catalogue holds no module of that name, and no build request the "+
"controller kept (the last 30 days) is for a directory of that name", inventory.ErrNoSuchModule, module)
if near := closestNames(ctx, inv, module); len(near) > 0 {
said += "; the closest names it holds: " + strings.Join(near, ", ")
}
return notHeld{kind: unknownModule, said: said}, nil
}
// notInCatalogue is the modules named that the catalogue does not hold.
func notInCatalogue(ctx context.Context, open *stores, modules []string) ([]string, error) {
shelf, err := open.inventory.Catalogue(ctx)
if err != nil {
return nil, err
}
var missing []string
for _, m := range modules {
if _, ok := shelf[m]; !ok {
missing = append(missing, m)
}
}
return missing, nil
}
// openPending is the open pending assignment of a module to a machine, if there is one.
func openPending(ctx context.Context, inv *inventory.Inventory, node, module string) (*inventory.PendingAssignment, error) {
rows, err := inv.PendingFor(ctx, node, module)
if err != nil {
return nil, err
}
for i := range rows {
if rows[i].Open() {
return &rows[i], nil
}
}
return nil, nil
}
// notRegistered answers an assignment of a module the catalogue does not hold: pending on a build in flight,
// known and not built, or unknown (novox/hq issue 325). An answer with no error is a pending assignment kept.
func notRegistered(ctx context.Context, open *stores, node, module string, opts assignOptions) (string, error) {
inv := open.inventory
if _, err := inv.NodeByName(ctx, node); err != nil {
return "", err
}
now := time.Now()
where, err := whereIs(ctx, inv, module, now)
if err != nil {
return "", err
}
waiting, err := openPending(ctx, inv, node, module)
if err != nil {
return "", err
}
switch where.kind {
case buildInFlight:
lead := where.said
if opts.Build {
lead += "\n no second build was asked: this one is running"
}
return keepPending(ctx, open, node, module, where.request, waiting, lead)
case knownNotBuilt:
last := where.request
if !opts.Build {
also := ""
if waiting != nil {
also = fmt.Sprintf("\n %s already has a pending assignment of %s, waiting for build %s; build "+
"\"true\" asks for a new build and makes it wait for that one", node, module, waiting.Build)
}
return "", fmt.Errorf("%w: %s%s\n assign it again with build \"true\" to ask for its build and keep this "+
"assignment until the build registers it; or `build` it (repository %s, path %s) and assign it "+
"once it is registered", inventory.ErrNoSuchModule, where.said, also, last.Repository, orRoot(last.Path))
}
if waiting != nil && waiting.State == inventory.PendingApplying {
return "", fmt.Errorf("%s is being assigned %s now; nothing was asked", node, module)
}
source := buildSource{Repository: last.Repository, Seat: last.Seat}
id, err := askABuild(ctx, source, last.Path, last.Ref)
if err != nil {
return "", fmt.Errorf("%s\n its build could not be asked for: %w; nothing was assigned", where.said, err)
}
asked := inventory.RequestOutcome{BuildRequest: inventory.BuildRequest{ID: id, Repository: last.Repository,
Seat: last.Seat, Path: last.Path, Ref: last.Ref, For: "assign", At: now.UTC()}}
recordBuildRequest(ctx, inv, asked.BuildRequest)
lead := fmt.Sprintf("%s\n build %s of %s is asked for now; %s can be assigned once it registers it",
where.said, id, requestSource(asked), module)
return keepPending(ctx, open, node, module, asked, waiting, lead)
}
answer := strings.TrimPrefix(where.said, inventory.ErrNoSuchModule.Error()) +
"\n a module in a repository is built with `build` (its repository and path), then assigned"
if opts.Build {
answer += "\n build \"true\" asks for nothing here: the controller has no build request saying where a " +
"module of that name is"
}
return "", fmt.Errorf("%w%s", inventory.ErrNoSuchModule, answer)
}
// keepPending records the assignment as pending on a build request — or, where one is already open, makes it
// wait for that build — and says so.
func keepPending(ctx context.Context, open *stores, node, module string, r inventory.RequestOutcome,
waiting *inventory.PendingAssignment, lead string) (string, error) {
inv := open.inventory
if waiting != nil {
if waiting.Build == r.ID {
return lead + fmt.Sprintf("\n %s already waits for %s on this build: nothing changed", node, module), nil
}
moved, err := inv.RepointPending(ctx, waiting.ID, r.ID, r.Repository, r.Path)
if err != nil {
return "", err
}
if !moved {
return lead + fmt.Sprintf("\n the pending assignment of %s to %s ended while this was asked: `status` "+
"says how", module, node), nil
}
return lead + fmt.Sprintf("\n the pending assignment of %s to %s, which waited for build %s, now waits "+
"for build %s", module, node, waiting.Build, r.ID), nil
}
_, err := inv.RecordPending(ctx, inventory.PendingAssignment{Node: node, Module: module, Build: r.ID,
Repository: r.Repository, Path: r.Path})
if errors.Is(err, inventory.ErrAlreadyPending) {
return lead + fmt.Sprintf("\n %s already waits for %s: nothing changed", node, module), nil
}
if err != nil {
return "", err
}
// The controller makes it, not this act (ADR 0261): when the build registers the module, or at its next
// tick if the module was registered between the look and the record.
return lead + fmt.Sprintf("\n the assignment is kept as pending: the controller assigns %s to %s when the "+
"build registers it, and `status` lists it until then. If the build fails it expires, says why and "+
"raises a condition; `unassign %s %s` withdraws it", module, node, node, module), nil
}
// inFlight is the newest build request for the module still to register it: no outcome yet within the
// bound, or a successful outcome heard a moment ago and being registered.
func inFlight(requests []inventory.RequestOutcome, now time.Time) (inventory.RequestOutcome, bool) {
for _, r := range requests {
if stillComing(r, now) {
return r, true
}
}
return inventory.RequestOutcome{}, false
}
func stillComing(r inventory.RequestOutcome, now time.Time) bool {
if r.Heard {
return r.Failed == "" && now.Sub(r.HeardAt) < registerGrace
}
return r.NotAsked == "" && now.Sub(r.At) < buildRequestBound
}
// whyNotBuilt says what came of a module's last build request. What was heard comes first: an outcome is the
// last word whatever its asker said.
func whyNotBuilt(r inventory.RequestOutcome, now time.Time) string {
where := fmt.Sprintf("%s in %s", orRoot(r.Path), r.Repository)
switch {
case r.Heard && r.Failed != "":
return fmt.Sprintf("%s; its last build, %s at %s, failed: %s", where, r.ID, clock(r.HeardAt),
firstLine(r.Failed))
case r.Heard && r.Module != "" && r.Module != r.Name():
return fmt.Sprintf("%s; its last build, %s, built it as %s", where, r.ID, r.Module)
case r.Heard:
return fmt.Sprintf("%s; its last build, %s at %s, was recorded and not registered — `builds` says why",
where, r.ID, clock(r.HeardAt))
case r.NotAsked != "" && r.For == "merge":
return fmt.Sprintf("%s; the merge that added it (%s) could not ask for its build: %s", where,
short(r.Commit), firstLine(r.NotAsked))
case r.NotAsked != "":
return fmt.Sprintf("%s; build %s, asked by %s, was not handed over: %s", where, r.ID,
askerOr(r.For), firstLine(r.NotAsked))
case r.OutcomeUnknown != "":
return fmt.Sprintf("%s; build %s was asked at %s, its asker stopped waiting (%s), and no outcome has "+
"been heard in %s", where, r.ID, clock(r.At), firstLine(r.OutcomeUnknown), now.Sub(r.At).Round(time.Minute))
default:
return fmt.Sprintf("%s; build %s was asked at %s, and no outcome has been heard in %s", where, r.ID,
clock(r.At), now.Sub(r.At).Round(time.Minute))
}
}
func askerOr(who string) string {
if who == "" {
return "the controller"
}
return who
}
// requestSource is a build request's source as a person reads it: repository@commit (directory).
func requestSource(r inventory.RequestOutcome) string {
at := short(r.Commit)
if at == "" {
at = r.Ref
}
if at == "" {
at = "its default branch"
}
return fmt.Sprintf("%s@%s (%s)", r.Repository, at, orRoot(r.Path))
}
// clock is a moment as a person reads it, in the controller's local time.
func clock(t time.Time) string { return t.Local().Format("2006-01-02 15:04:05 MST") }
// closestNames is up to three names near the one asked: registered modules and directories the controller
// asked to build.
func closestNames(ctx context.Context, inv *inventory.Inventory, name string) []string {
var candidates []string
if shelf, err := inv.Catalogue(ctx); err == nil {
for m := range shelf {
candidates = append(candidates, m)
}
}
if asked, err := inv.RequestedNames(ctx); err == nil {
candidates = append(candidates, asked...)
}
type scored struct {
name string
distance int
}
seen := map[string]bool{}
var near []scored
limit := max(2, len(name)/3)
for _, c := range candidates {
if seen[c] || c == name {
continue
}
seen[c] = true
d := editDistance(name, c)
if d <= limit || strings.Contains(c, name) || strings.Contains(name, c) {
near = append(near, scored{c, d})
}
}
sort.Slice(near, func(i, j int) bool {
if near[i].distance != near[j].distance {
return near[i].distance < near[j].distance
}
return near[i].name < near[j].name
})
var out []string
for i := 0; i < len(near) && i < 3; i++ {
out = append(out, near[i].name)
}
return out
}
// editDistance is the Levenshtein distance between two names.
func editDistance(a, b string) int {
ra, rb := []rune(a), []rune(b)
prev := make([]int, len(rb)+1)
for j := range prev {
prev[j] = j
}
for i := 1; i <= len(ra); i++ {
cur := make([]int, len(rb)+1)
cur[0] = i
for j := 1; j <= len(rb); j++ {
cost := 1
if ra[i-1] == rb[j-1] {
cost = 0
}
cur[j] = min(prev[j]+1, cur[j-1]+1, prev[j-1]+cost)
}
prev = cur
}
return prev[len(rb)]
}
// applyPending makes a pending assignment now that its module is registered (ADR 0261). Under the machine's
// hold it claims the row (waiting → applying), so a withdrawal cannot land between the look and the act,
// then assigns by the same act a person's assign is, then says what came of it: applied, or refused with the
// refusal. Returns what it said, or nothing when the row no longer waited.
func applyPending(ctx context.Context, open *stores, p inventory.PendingAssignment, by string) string {
inv := open.inventory
ctx, release, err := holdNodes(ctx, open, []string{p.Node})
if err != nil {
// Left waiting: the next tick tries again.
return fmt.Sprintf("the pending assignment of %s to %s waits: %s could not be held: %v", p.Module, p.Node,
p.Node, err)
}
defer release()
claimed, err := inv.ClaimPending(ctx, p.ID)
if err != nil {
return fmt.Sprintf("the pending assignment of %s to %s could not be taken for making: %v", p.Module, p.Node, err)
}
if !claimed {
return ""
}
answer, err := assign(ctx, open, p.Node, p.Module)
state, note := inventory.PendingApplied, ""
switch {
case err != nil:
state = inventory.PendingRefused
note = fmt.Sprintf("the pending assignment of %s to %s was refused when build %s registered it: %s",
p.Module, p.Node, by, firstLine(err.Error()))
case strings.Contains(answer, "already runs"):
note = fmt.Sprintf("%s already runs %s: the pending assignment is met", p.Node, p.Module)
default:
// The same as a person's assign: nothing is sent by this act.
note = fmt.Sprintf("%s is assigned %s: build %s registered it (pending since %s); `push %s` sends it",
p.Node, p.Module, by, clock(p.Since), p.Node)
}
settled, serr := inv.SettlePending(ctx, p.ID, inventory.PendingApplying, state, note)
if serr != nil {
return fmt.Sprintf("%s — and it could not be recorded as settled: %v", note, serr)
}
if !settled {
return ""
}
return note
}
// expirePending ends a waiting pending assignment whose build will not register its module, with why.
func expirePending(ctx context.Context, inv *inventory.Inventory, p inventory.PendingAssignment, why string) string {
note := fmt.Sprintf("the pending assignment of %s to %s expired: %s; nothing was assigned. Assign it again with "+
"build \"true\" to ask for the build again", p.Module, p.Node, why)
settled, err := inv.SettlePending(ctx, p.ID, inventory.PendingWaiting, inventory.PendingExpired, note)
if err != nil {
return fmt.Sprintf("%s — and it could not be recorded: %v", note, err)
}
if !settled {
return ""
}
return note
}
// settlePendingOnBuild is a build's outcome reaching the assignments waiting for it, wherever the outcome is
// taken in (novox/hq issue 325): a module registered makes every assignment pending on it; a build of a
// pending assignment that failed, or was not registered, ends it with why. Returns what was said.
func settlePendingOnBuild(ctx context.Context, open *stores, result link.BuildResult, module string, takeErr error) []string {
inv := open.inventory
waiting, err := inv.Pending(ctx, time.Time{})
if err != nil {
return []string{fmt.Sprintf("the pending assignments could not be read after build %s: %v", result.ID, err)}
}
registered := module != "" && (takeErr == nil || errors.Is(takeErr, inventory.ErrSuperseded))
var said []string
for _, p := range waiting {
switch {
case registered && p.Module == module:
if line := applyPending(ctx, open, p, result.ID); line != "" {
said = append(said, line)
}
case !registered && p.Build == result.ID:
why := firstLine(result.Failed)
if why == "" && takeErr != nil {
why = firstLine(takeErr.Error())
}
what := fmt.Sprintf("build %s of %s failed: %s", result.ID, orRoot(p.Path), why)
if result.Failed == "" {
what = fmt.Sprintf("build %s of %s was recorded and not registered: %s", result.ID, orRoot(p.Path), why)
}
if line := expirePending(ctx, inv, p, what); line != "" {
said = append(said, line)
}
}
}
return said
}
// settlingPending is the controller's tick over the pending assignments (ADR 0261), every settleEvery until
// the context ends. What it does is printed to the controller's log, kept in each row's note (which `status`
// lists for a day), and raised as a condition for an assignment that was not made.
func settlingPending(ctx context.Context, open *stores) {
for {
for _, line := range settlePending(ctx, open, time.Now()) {
fmt.Println(line)
}
select {
case <-ctx.Done():
return
case <-time.After(settleEvery):
}
}
}
// settlePending is one tick: claims left by a controller that stopped are settled from what the mesh holds;
// a module registered meanwhile is assigned; a build heard and not registered, or silent past its bound,
// ends its assignment; ended ones are raised as conditions and cleared once answered; ended rows older than
// KeptFor are deleted. Returns what it did.
func settlePending(ctx context.Context, open *stores, now time.Time) []string {
inv := open.inventory
var said []string
say := func(format string, args ...any) { said = append(said, fmt.Sprintf(format, args...)) }
shelf, err := inv.Catalogue(ctx)
if err != nil {
say("the pending assignments are not settled this time: the catalogue could not be read: %v", err)
return said
}
stuck, err := inv.Stuck(ctx, now.Add(-claimStaleAfter))
if err != nil {
say("claims left by a stopped controller could not be read: %v", err)
}
for _, p := range stuck {
assigned, err := inv.Assigned(ctx, p.Node)
if err != nil {
say("the claim on %s for %s could not be settled: %v", p.Node, p.Module, err)
continue
}
if containsString(assigned, p.Module) {
if _, err := inv.SettlePending(ctx, p.ID, inventory.PendingApplying, inventory.PendingApplied,
fmt.Sprintf("%s is assigned %s (settled after the controller that made it stopped)", p.Node, p.Module)); err != nil {
say("the claim on %s for %s could not be settled: %v", p.Node, p.Module, err)
}
continue
}
if err := inv.ReleaseClaim(ctx, p.ID); err != nil {
say("the claim on %s for %s could not be released: %v", p.Node, p.Module, err)
}
}
waiting, err := inv.Pending(ctx, time.Time{})
if err != nil {
say("the pending assignments could not be read to settle them: %v", err)
}
for _, p := range waiting {
if _, ok := shelf[p.Module]; ok {
if line := applyPending(ctx, open, p, p.Build); line != "" {
said = append(said, line)
}
continue
}
requests, err := inv.RequestsNamed(ctx, p.Module)
if err != nil {
say("the build requests for %s could not be read: %v", p.Module, err)
continue
}
var r *inventory.RequestOutcome
for i := range requests {
if requests[i].ID == p.Build {
r = &requests[i]
}
}
why := ""
switch {
case r == nil && now.Sub(p.Since) > buildRequestBound:
why = fmt.Sprintf("build %s is no longer on record", p.Build)
case r != nil && !stillComing(*r, now) && (r.Heard || r.NotAsked != "" || r.OutcomeUnknown != ""):
why = whyNotBuilt(*r, now)
case r != nil && !stillComing(*r, now):
why = fmt.Sprintf("no outcome of build %s was heard within %s of asking", p.Build, buildRequestBound)
}
if why == "" {
continue
}
if line := expirePending(ctx, inv, p, why); line != "" {
said = append(said, line)
}
}
said = append(said, raisePendingEnded(ctx, inv)...)
if n, err := inv.ForgetEndedPending(ctx, now); err != nil {
say("ended pending assignments could not be deleted: %v", err)
} else if n > 0 {
say("deleted %d pending assignment(s) ended more than %s ago", n, inventory.KeptFor)
}
return said
}
func containsString(xs []string, x string) bool {
for _, y := range xs {
if y == x {
return true
}
}
return false
}
// pendingObservation is the condition for a pending assignment that was not made: one per pending
// assignment, its row's id the last part of its id, so one row's condition is never another's.
func pendingObservation(p inventory.PendingAssignment) conditions.Observation {
return conditions.Observation{Scope: conditions.ScopeMachine, ID: fmt.Sprintf("%s.%s.%d", p.Node, p.Module, p.ID),
Token: kindPendingEnded, Kind: kindPendingEnded, Machine: p.Node, Severity: conditions.Warning,
Resolver: conditions.ResolverOperator, Source: "pending assignments",
Summary: p.Note}
}
// raisePendingEnded raises a condition for each pending assignment that expired or was refused, once, and
// clears it when it is answered: the module was assigned to that machine since, a newer pending assignment of
// it is open, or a person took it back with `unassign`. Where this process keeps no conditions, nothing is
// stamped, and the serving controller's tick raises it.
func raisePendingEnded(ctx context.Context, inv *inventory.Inventory) []string {
keeper := conditionsFrom
if keeper == nil {
return nil
}
var said []string
toRaise, err := inv.ToRaise(ctx)
if err != nil {
return []string{fmt.Sprintf("pending assignments not made could not be read to raise them: %v", err)}
}
for _, p := range toRaise {
if _, err := keeper.Observe(ctx, pendingObservation(p)); err != nil {
said = append(said, fmt.Sprintf("the condition for %s on %s could not be raised: %v", p.Module, p.Node, err))
continue
}
if err := inv.MarkPending(ctx, p.ID, "raised"); err != nil {
said = append(said, fmt.Sprintf("the condition for %s on %s was raised and not recorded: %v", p.Module, p.Node, err))
}
}
toClear, err := inv.ToClear(ctx)
if err != nil {
return append(said, fmt.Sprintf("pending assignments raised could not be read to clear them: %v", err))
}
for _, p := range toClear {
why, answered, err := pendingAnswered(ctx, inv, p)
if err != nil {
said = append(said, fmt.Sprintf("whether %s on %s is answered could not be read: %v", p.Module, p.Node, err))
continue
}
if !answered {
continue
}
if _, err := keeper.Clear(ctx, pendingObservation(p).Key(), why); err != nil {
said = append(said, fmt.Sprintf("the condition for %s on %s could not be cleared: %v", p.Module, p.Node, err))
continue
}
if err := inv.MarkPending(ctx, p.ID, "cleared"); err != nil {
said = append(said, fmt.Sprintf("the condition for %s on %s was cleared and not recorded: %v", p.Module, p.Node, err))
}
}
return append(said, clearOrphaned(ctx, keeper, inv)...)
}
// pendingAnswered says whether a pending assignment that was not made has been answered since, and how.
func pendingAnswered(ctx context.Context, inv *inventory.Inventory, p inventory.PendingAssignment) (string, bool, error) {
if p.Acknowledged != nil {
return "taken back with unassign", true, nil
}
assigned, err := inv.Assigned(ctx, p.Node)
if err != nil {
return "", false, err
}
if containsString(assigned, p.Module) {
return p.Module + " is assigned to " + p.Node + " since", true, nil
}
rows, err := inv.PendingFor(ctx, p.Node, p.Module)
if err != nil {
return "", false, err
}
// A newer pending assignment answers it only while it is open or once it was made: one that itself
// ended unmade is its own condition, and answers nothing.
for _, r := range rows {
if r.ID != p.ID && r.Since.After(p.Since) && (r.Open() || r.State == inventory.PendingApplied) {
return "assigned again, pending on another build", true, nil
}
}
return "", false, nil
}
// clearOrphaned clears every open condition of an assignment not made whose pending assignment is no longer
// on record: its machine was removed, which takes its pending assignments with it.
func clearOrphaned(ctx context.Context, keeper *conditions.Keeper, inv *inventory.Inventory) []string {
open, err := keeper.Open(ctx)
if err != nil {
return []string{fmt.Sprintf("the open conditions could not be read to clear orphaned ones: %v", err)}
}
byID := map[int64]string{}
var ids []int64
for _, c := range open {
if c.Kind != kindPendingEnded {
continue
}
// `<scope>.<node>.<module>.<row>.<kind>`: the row is the part before the kind.
parts := strings.Split(c.Key, ".")
if len(parts) < 2 {
continue
}
var id int64
if _, err := fmt.Sscan(parts[len(parts)-2], &id); err != nil {
continue
}
byID[id] = c.Key
ids = append(ids, id)
}
if len(ids) == 0 {
return nil
}
known, err := inv.PendingKnown(ctx, ids)
if err != nil {
return []string{fmt.Sprintf("the pending assignments could not be read to clear orphaned conditions: %v", err)}
}
var said []string
for id, key := range byID {
if known[id] {
continue
}
if _, err := keeper.Clear(ctx, key, "its pending assignment is no longer on record: its machine was removed"); err != nil {
said = append(said, fmt.Sprintf("the condition %s could not be cleared: %v", key, err))
}
}
return said
}
// withdrawPending is `unassign` of a module only pending on a machine, under the machine's hold: a waiting
// pending assignment is withdrawn; one being made now is refused, never overridden; one that expired or was
// refused is taken back, which answers its condition. False when there is none.
func withdrawPending(ctx context.Context, inv *inventory.Inventory, node, module string) (string, bool, error) {
rows, err := inv.PendingFor(ctx, node, module)
if err != nil {
return "", false, err
}
var takenBack []string
for _, p := range rows {
switch p.State {
case inventory.PendingApplying:
return "", true, fmt.Errorf("%s is being assigned %s now, as build %s registered it: nothing was "+
"withdrawn; `unassign %s %s` again once it is assigned takes it off", node, module, p.Build, node, module)
case inventory.PendingWaiting:
note := fmt.Sprintf("the pending assignment of %s to %s is withdrawn; build %s goes on, and registers "+
"%s assigned nowhere", module, node, p.Build, module)
settled, err := inv.SettlePending(ctx, p.ID, inventory.PendingWaiting, inventory.PendingWithdrawn, note)
if err != nil {
return "", false, err
}
if !settled {
return "", true, fmt.Errorf("the pending assignment of %s to %s changed while it was withdrawn: "+
"`status` says how; nothing was withdrawn", module, node)
}
return note, true, nil
case inventory.PendingExpired, inventory.PendingRefused:
if p.Cleared != nil || p.Acknowledged != nil {
continue
}
if err := inv.MarkPending(ctx, p.ID, "acknowledged"); err != nil {
return "", false, err
}
takenBack = append(takenBack, fmt.Sprintf("build %s (%s)", p.Build, p.State))
}
}
if len(takenBack) > 0 {
return fmt.Sprintf("the pending assignment(s) of %s to %s that were not made are taken back: %s; their "+
"conditions clear", module, node, strings.Join(takenBack, ", ")), true, nil
}
return "", false, nil
}
+864
View File
@@ -0,0 +1,864 @@
package main
import (
"context"
"encoding/json"
"errors"
"fmt"
"slices"
"strings"
"testing"
"time"
"github.com/novox/mesh-controller/internal/catalogue"
"github.com/novox/mesh-controller/internal/conditions"
"github.com/novox/mesh-controller/internal/inventory"
"github.com/novox/mesh-controller/internal/link"
)
// failedBuild settles the merge's build of modules/sensors as failed.
func failedBuild(t *testing.T, open *stores, id string) {
t.Helper()
if err := (builds{open.inventory, open}).Built(t.Context(), outcome(id, "modules/sensors", "3da80a4b00aa",
"", "boom")); err != nil {
t.Fatal(err)
}
}
func conditionOpen(t *testing.T, p inventory.PendingAssignment) bool {
t.Helper()
_, found, err := conditionsFrom.Get(t.Context(), pendingObservation(p).Key())
if err != nil {
t.Fatal(err)
}
return found
}
func rowOf(t *testing.T, open *stores, id int64) inventory.PendingAssignment {
t.Helper()
rows, err := open.inventory.Pending(t.Context(), time.Unix(0, 0))
if err != nil {
t.Fatal(err)
}
for _, r := range rows {
if r.ID == id {
return r
}
}
t.Fatalf("pending assignment %d is not on record", id)
return inventory.PendingAssignment{}
}
// novox/hq issue 325: an assignment of a module the catalogue does not hold says which case it is in — a
// build in flight (kept pending), known and not built (said, with build), or unknown (refused, with the
// closest names) — and a pending assignment is made when its build registers the module, or ends with why.
// aCatalogueMesh is aMesh with one module held from the catalogue, so a merge of it is acted on.
func aCatalogueMesh(t *testing.T) *stores {
t.Helper()
open := aMesh(t)
if err := open.inventory.RegisterModule(t.Context(), catalogue.Manifest{Module: "networkmanager", Version: "1"},
inventory.Source{Repository: "novox/mesh-catalog", Seat: "git", Path: "modules/networkmanager", Ref: "main",
BuiltFrom: "c0", Head: "c0"}); err != nil {
t.Fatal(err)
}
return open
}
// mergeAdding is the catalogue's merge adding one module's directory, acted on as the bus hands it over.
func mergeAdding(t *testing.T, open *stores, dir, commit string) {
t.Helper()
m := link.SourceMoved{Owner: "novox", Repo: "mesh-catalog", Base: "main", Commit: commit,
Paths: []string{dir + "/module.json"}, ModuleDirs: []string{dir}, ModuleDirsSaid: true}
if err := (following{open: open}).SourceMoved(t.Context(), m); err != nil {
t.Fatal(err)
}
}
// outcome is a build's outcome as the build seat announces it: registered as module, or failed.
func outcome(id, dir, commit, module, failed string) link.BuildResult {
r := link.BuildResult{ID: id, Repository: "novox/mesh-catalog", Path: dir, Ref: "main", Commit: commit,
On: "anchor", Failed: failed, Source: &link.SourceOnSeat{Repository: "novox/mesh-catalog", Seat: "git"}}
if failed == "" {
r.Module = module
r.Manifest, _ = json.Marshal(catalogue.Manifest{Module: module, Version: "1"})
}
return r
}
func assignedTo(t *testing.T, open *stores, node string) []string {
t.Helper()
got, err := open.inventory.Assigned(t.Context(), node)
if err != nil {
t.Fatal(err)
}
return got
}
func pendingOf(t *testing.T, open *stores) []inventory.PendingAssignment {
t.Helper()
got, err := open.inventory.Pending(t.Context(), time.Now().Add(-time.Hour))
if err != nil {
t.Fatal(err)
}
return got
}
// A build in flight: the merge asked for it, the request is kept as the merge's, the assignment is kept
// pending and said so with the build, status reads it without settling anything, and the build's outcome
// makes it — saying that a push sends it, and nothing more.
func TestAnAssignmentWaitsForTheBuildInFlight(t *testing.T) {
open := aCatalogueMesh(t)
ctx := t.Context()
asksWithPaths(t)
mergeAdding(t, open, "modules/sensors", "3da80a4b00aa")
requests, err := open.inventory.RequestsNamed(ctx, "sensors")
if err != nil || len(requests) != 1 || requests[0].For != "merge" || requests[0].Commit != "3da80a4b00aa" {
t.Fatalf("the merge's build request: %+v %v", requests, err)
}
said, err := assign(ctx, open, "laptop", "sensors")
if err != nil {
t.Fatalf("an assignment while its build runs was refused: %v", err)
}
t.Logf("assign laptop sensors, while its build runs:\n%s", said)
for _, want := range []string{"being built from novox/mesh-catalog@3da80a4b (modules/sensors)",
"build b-modules/sensors", "kept as pending", "the controller assigns sensors to laptop", "`unassign laptop sensors`"} {
if !strings.Contains(said, want) {
t.Fatalf("the answer does not say %q:\n%s", want, said)
}
}
if slices.Contains(assignedTo(t, open, "laptop"), "sensors") {
t.Fatal("a module not registered was assigned")
}
pending := pendingOf(t, open)
if len(pending) != 1 || pending[0].State != inventory.PendingWaiting || pending[0].Build != "b-modules/sensors" {
t.Fatalf("pending: %+v", pending)
}
if again, err := assign(ctx, open, "laptop", "sensors"); err != nil || !strings.Contains(again, "already waits") {
t.Fatalf("a second assignment: %q %v", again, err)
}
asked, err := theThreeQuestions(ctx, open)
if err != nil {
t.Fatal(err)
}
if asked.well() {
t.Fatal("status called the mesh well while an assignment waits")
}
body, err := statusAsJSON(asked)
if err != nil {
t.Fatal(err)
}
var doc meshStatus
if err := json.Unmarshal(body, &doc); err != nil {
t.Fatal(err)
}
if len(doc.Pending) != 1 || doc.Pending[0].State != "waiting" || doc.Pending[0].Module != "sensors" {
t.Fatalf("status says pending %+v", doc.Pending)
}
if err := (builds{open.inventory, open}).Built(ctx, outcome("b-modules/sensors", "modules/sensors",
"3da80a4b00aa", "sensors", "")); err != nil {
t.Fatal(err)
}
if !slices.Contains(assignedTo(t, open, "laptop"), "sensors") {
t.Fatal("the build registered sensors and the pending assignment was not made")
}
pending = pendingOf(t, open)
if len(pending) != 1 || pending[0].State != inventory.PendingApplied ||
!strings.HasSuffix(pending[0].Note, "`push laptop` sends it") {
t.Fatalf("pending after the build: %+v", pending)
}
}
// **A read settles nothing** (review of #150, point 6): status — and so the board, the summary and the
// probe, which compose from the same reading — leaves a pending assignment whose module is registered as it
// is; the controller's tick makes it.
func TestAStatusReadSettlesNothing(t *testing.T) {
open := aCatalogueMesh(t)
ctx := t.Context()
asksWithPaths(t)
mergeAdding(t, open, "modules/sensors", "3da80a4b00aa")
if _, err := assign(ctx, open, "laptop", "sensors"); err != nil {
t.Fatal(err)
}
register(t, open, catalogue.Manifest{Module: "sensors", Version: "1"})
for range 2 {
if _, err := theThreeQuestions(ctx, open); err != nil {
t.Fatal(err)
}
}
if slices.Contains(assignedTo(t, open, "laptop"), "sensors") || pendingOf(t, open)[0].State != inventory.PendingWaiting {
t.Fatalf("a status read settled a pending assignment: %+v", pendingOf(t, open))
}
said := settlePending(ctx, open, time.Now())
if !slices.Contains(assignedTo(t, open, "laptop"), "sensors") || pendingOf(t, open)[0].State != inventory.PendingApplied {
t.Fatalf("the tick did not make it: %v %+v", said, pendingOf(t, open))
}
}
// The build of a pending assignment fails: it expires with the build's words and nothing is assigned; the
// tick raises it as a condition once; it then reads as known and not built; and `unassign` takes it back,
// after which the tick clears the condition. Status is not called well while the condition is open, and is
// again once it clears — no fixed day of "not well" (review point 6).
func TestAPendingAssignmentExpiresWhenItsBuildFails(t *testing.T) {
open := aCatalogueMesh(t)
ctx := t.Context()
asksWithPaths(t)
mergeAdding(t, open, "modules/sensors", "3da80a4b00aa")
if _, err := assign(ctx, open, "laptop", "sensors"); err != nil {
t.Fatal(err)
}
if err := (builds{open.inventory, open}).Built(ctx, outcome("b-modules/sensors", "modules/sensors",
"3da80a4b00aa", "", "go build: undefined: sensorsRead")); err != nil {
t.Fatal(err)
}
pending := pendingOf(t, open)
if len(pending) != 1 || pending[0].State != inventory.PendingExpired {
t.Fatalf("pending after a failed build: %+v", pending)
}
for _, want := range []string{"expired", "build b-modules/sensors of modules/sensors failed",
"undefined: sensorsRead", "nothing was assigned", "build \"true\""} {
if !strings.Contains(pending[0].Note, want) {
t.Fatalf("the expiry does not say %q: %s", want, pending[0].Note)
}
}
if slices.Contains(assignedTo(t, open, "laptop"), "sensors") {
t.Fatal("a failed build's module was assigned")
}
settlePending(ctx, open, time.Now())
key := pendingObservation(pending[0]).Key()
c, found, err := conditionsFrom.Get(ctx, key)
if err != nil || !found || c.Kind != kindPendingEnded {
t.Fatalf("no condition %s for the assignment not made: %+v %v %v", key, c, found, err)
}
if asked, err := theThreeQuestions(ctx, open); err != nil || asked.well() || len(asked.conditions) == 0 {
t.Fatalf("status does not hold the open condition: %v", err)
}
_, err = assign(ctx, open, "laptop", "sensors")
if !errors.Is(err, inventory.ErrNoSuchModule) {
t.Fatalf("a module whose build failed: %v", err)
}
for _, want := range []string{"sensors is known and not registered", "its last build, b-modules/sensors",
"failed: go build: undefined: sensorsRead", "build \"true\"", "repository novox/mesh-catalog, path modules/sensors"} {
if !strings.Contains(err.Error(), want) {
t.Fatalf("the refusal does not say %q:\n%v", want, err)
}
}
said, err := unassign(ctx, open, "laptop", "sensors")
if err != nil || !strings.Contains(said, "taken back") {
t.Fatalf("unassign of an expired pending assignment: %q %v", said, err)
}
settlePending(ctx, open, time.Now())
if _, found, _ := conditionsFrom.Get(ctx, key); found {
t.Fatal("the condition stayed open after the assignment was taken back")
}
// Nothing of it keeps status from being well any more: no open pending assignment, no open condition of it.
asked, err := theThreeQuestions(ctx, open)
if err != nil || pendingOpen(asked.pending) {
t.Fatalf("status still counts the pending assignment: %+v %v", asked.pending, err)
}
for _, c := range asked.conditions {
if c.Kind == kindPendingEnded {
t.Fatalf("status still holds the condition: %+v", c)
}
}
}
// The condition also clears when the module is later assigned to the machine.
func TestTheConditionClearsWhenTheModuleIsAssignedLater(t *testing.T) {
open := aCatalogueMesh(t)
ctx := t.Context()
asksWithPaths(t)
mergeAdding(t, open, "modules/sensors", "3da80a4b00aa")
if _, err := assign(ctx, open, "laptop", "sensors"); err != nil {
t.Fatal(err)
}
if err := (builds{open.inventory, open}).Built(ctx, outcome("b-modules/sensors", "modules/sensors",
"3da80a4b00aa", "", "boom")); err != nil {
t.Fatal(err)
}
settlePending(ctx, open, time.Now())
key := pendingObservation(pendingOf(t, open)[0]).Key()
register(t, open, catalogue.Manifest{Module: "sensors", Version: "1"})
if _, err := assign(ctx, open, "laptop", "sensors"); err != nil {
t.Fatal(err)
}
settlePending(ctx, open, time.Now())
if _, found, _ := conditionsFrom.Get(ctx, key); found {
t.Fatal("the condition stayed open after the module was assigned")
}
}
// Known and not built, asked with build: the build is asked from where the last one was, the request kept as
// assign's, and the assignment kept pending on the new build.
func TestAssignWithBuildAsksForAKnownModule(t *testing.T) {
open := aCatalogueMesh(t)
ctx := t.Context()
inv := open.inventory
if err := inv.RecordBuildRequest(ctx, inventory.BuildRequest{ID: "build-old", Repository: "novox/mesh-catalog",
Seat: "git", Path: "modules/sensors", Ref: "main", For: "build", At: time.Now().Add(-time.Hour)}); err != nil {
t.Fatal(err)
}
if err := inv.RecordBuild(ctx, inventory.Build{ID: "build-old", Repository: "novox/mesh-catalog", Ref: "main",
Path: "modules/sensors", On: "anchor", Failed: "no space left on device"}); err != nil {
t.Fatal(err)
}
asked := asksWithPaths(t)
said, err := assignWith(ctx, open, "laptop", assignOptions{Build: true}, "sensors")
if err != nil {
t.Fatal(err)
}
if len(*asked) != 1 || (*asked)[0] != [3]string{"novox/mesh-catalog", "modules/sensors", "main"} {
t.Fatalf("asked %v", *asked)
}
for _, want := range []string{"known and not registered", "no space left on device",
"build b-modules/sensors of novox/mesh-catalog@main (modules/sensors) is asked for now", "kept as pending"} {
if !strings.Contains(said, want) {
t.Fatalf("the answer does not say %q:\n%s", want, said)
}
}
pending := pendingOf(t, open)
if len(pending) != 1 || pending[0].Build != "b-modules/sensors" {
t.Fatalf("pending: %+v", pending)
}
requests, _ := inv.RequestsNamed(ctx, "sensors")
if len(requests) != 2 || requests[0].ID != "b-modules/sensors" || requests[0].For != "assign" {
t.Fatalf("the request is not kept as assign's: %+v", requests)
}
}
// **build "true" with a pending assignment already waiting** (review point 5): the waiting one is made to
// wait for the new build, and no build is asked that nothing waits on.
func TestAssignWithBuildRepointsTheWaitingAssignment(t *testing.T) {
open := aCatalogueMesh(t)
ctx := t.Context()
inv := open.inventory
if err := inv.RecordBuildRequest(ctx, inventory.BuildRequest{ID: "build-old", Repository: "novox/mesh-catalog",
Seat: "git", Path: "modules/sensors", Ref: "main", For: "merge", At: time.Now().Add(-time.Hour)}); err != nil {
t.Fatal(err)
}
if err := inv.RecordBuild(ctx, inventory.Build{ID: "build-old", Repository: "novox/mesh-catalog", Ref: "main",
Path: "modules/sensors", On: "anchor", Failed: "boom"}); err != nil {
t.Fatal(err)
}
// Waiting on the failed build, its expiry not yet settled.
if _, err := inv.RecordPending(ctx, inventory.PendingAssignment{Node: "laptop", Module: "sensors",
Build: "build-old", Repository: "novox/mesh-catalog", Path: "modules/sensors"}); err != nil {
t.Fatal(err)
}
_, err := assign(ctx, open, "laptop", "sensors")
if err == nil || !strings.Contains(err.Error(), "already has a pending assignment of sensors, waiting for build build-old") {
t.Fatalf("known and not built with a pending assignment waiting: %v", err)
}
asked := asksWithPaths(t)
said, err := assignWith(ctx, open, "laptop", assignOptions{Build: true}, "sensors")
if err != nil || !strings.Contains(said, "which waited for build build-old, now waits for build b-modules/sensors") {
t.Fatalf("assign with build: %q %v", said, err)
}
if len(*asked) != 1 {
t.Fatalf("asked %v", *asked)
}
pending := pendingOf(t, open)
if len(pending) != 1 || pending[0].Build != "b-modules/sensors" || pending[0].State != inventory.PendingWaiting {
t.Fatalf("pending: %+v", pending)
}
// The new build in flight now: build "true" again asks nothing.
if said, err := assignWith(ctx, open, "laptop", assignOptions{Build: true}, "sensors"); err != nil ||
!strings.Contains(said, "no second build was asked") || len(*asked) != 1 {
t.Fatalf("a second build true while the build runs: %q %v, asked %v", said, err, *asked)
}
}
// **A failed ask never reads as in flight** (review point 2): a merge whose ask could not be made keeps the
// request as not asked, and assign says the merge could not ask; a request whose asker could not hand it over
// or stopped waiting reads the same, unless its outcome was heard.
func TestAFailedAskIsNotABuildInFlight(t *testing.T) {
open := aCatalogueMesh(t)
ctx := t.Context()
inv := open.inventory
was := askABuild
askABuild = func(context.Context, buildSource, string, string) (string, error) {
return "", errors.New("cannot submit a build: nats: timeout")
}
t.Cleanup(func() { askABuild = was })
mergeAdding(t, open, "modules/sensors", "3da80a4b00aa")
_, err := assign(ctx, open, "laptop", "sensors")
if err == nil || !strings.Contains(err.Error(), "the merge that added it (3da80a4b) could not ask for its build: cannot submit") {
t.Fatalf("a merge that could not ask: %v", err)
}
if len(pendingOf(t, open)) != 0 {
t.Fatal("a pending assignment waits on a build never asked")
}
if err := inv.RecordBuildRequest(ctx, inventory.BuildRequest{ID: "build-waited", Repository: "novox/mesh-catalog",
Seat: "git", Path: "modules/gauges", Ref: "main", For: "build"}); err != nil {
t.Fatal(err)
}
if err := inv.MarkNotAsked(ctx, "build-waited", "cannot submit a build: no responders"); err != nil {
t.Fatal(err)
}
_, err = assign(ctx, open, "laptop", "gauges")
if err == nil || !strings.Contains(err.Error(), "build build-waited, asked by build, was not handed over") {
t.Fatalf("a build not handed over: %v", err)
}
// Heard first, the outcome stands: marking it afterwards changes nothing.
if err := inv.RecordBuild(ctx, inventory.Build{ID: "build-waited", Repository: "novox/mesh-catalog", Ref: "main",
Path: "modules/gauges", On: "anchor", Failed: "the real failure"}); err != nil {
t.Fatal(err)
}
_, err = assign(ctx, open, "laptop", "gauges")
if err == nil || !strings.Contains(err.Error(), "failed: the real failure") {
t.Fatalf("an outcome heard after a failed wait: %v", err)
}
}
// A plan's tier keeps its requests as the plan's (review point 3).
func TestAPlansBuildRequestIsThePlans(t *testing.T) {
open := aCatalogueMesh(t)
ctx := t.Context()
asksWithPaths(t)
if err := open.inventory.RegisterModule(ctx, catalogue.Manifest{Module: "app", Version: "1"},
inventory.Source{Repository: "novox/mesh-catalog", Seat: "git", Path: "modules/app", Ref: "main",
BuiltFrom: "c0", Head: "c0"}); err != nil {
t.Fatal(err)
}
m := link.SourceMoved{Owner: "novox", Repo: "mesh-catalog", Base: "main", Commit: "c1aaaaaaaa",
Paths: []string{"modules/app/index.ts"}, ModuleDirs: []string{"modules/app"}, ModuleDirsSaid: true}
if err := (following{open: open}).SourceMoved(ctx, m); err != nil {
t.Fatal(err)
}
requests, err := open.inventory.RequestsNamed(ctx, "app")
if err != nil || len(requests) != 1 || requests[0].For != "plan" {
t.Fatalf("the plan's request: %+v %v", requests, err)
}
}
// **A build heard as built and not registered yet is in flight** (review point 4): the outcome is recorded
// before the module is registered; for that moment assign keeps the assignment pending. Past the grace it is
// recorded and not registered, and the tick ends the pending assignment with that.
func TestABuildBeingRegisteredIsInFlight(t *testing.T) {
open := aCatalogueMesh(t)
ctx := t.Context()
inv := open.inventory
if err := inv.RecordBuildRequest(ctx, inventory.BuildRequest{ID: "build-done", Repository: "novox/mesh-catalog",
Seat: "git", Path: "modules/sensors", Ref: "main", For: "merge"}); err != nil {
t.Fatal(err)
}
if err := inv.RecordBuild(ctx, inventory.Build{ID: "build-done", Repository: "novox/mesh-catalog", Ref: "main",
Path: "modules/sensors", Module: "sensors", On: "anchor"}); err != nil {
t.Fatal(err)
}
said, err := assign(ctx, open, "laptop", "sensors")
if err != nil || !strings.Contains(said, "built and being registered") || !strings.Contains(said, "kept as pending") {
t.Fatalf("a build being registered: %q %v", said, err)
}
settlePending(ctx, open, time.Now().Add(registerGrace+time.Minute))
p := pendingOf(t, open)
if len(p) != 1 || p[0].State != inventory.PendingExpired || !strings.Contains(p[0].Note, "recorded and not registered") {
t.Fatalf("past the grace: %+v", p)
}
}
// Unknown: refused as no module of that name, claiming only what was looked at, with the closest names; build
// asks for nothing.
func TestAnUnknownModuleIsRefusedWithTheClosestNames(t *testing.T) {
open := aCatalogueMesh(t)
ctx := t.Context()
register(t, open, catalogue.Manifest{Module: "sensord", Version: "1"})
asked := asksWithPaths(t)
_, err := assignWith(ctx, open, "laptop", assignOptions{Build: true}, "sensors")
if !errors.Is(err, inventory.ErrNoSuchModule) {
t.Fatalf("an unknown module: %v", err)
}
for _, want := range []string{"no module of that name: sensors", "the catalogue holds no module of that name",
"no build request the controller kept (the last 30 days)", "the closest names it holds: sensord",
"asks for nothing here"} {
if !strings.Contains(err.Error(), want) {
t.Fatalf("the refusal does not say %q:\n%v", want, err)
}
}
if strings.Contains(err.Error(), "merge") {
t.Fatalf("the refusal claims a merge it never looked at: %v", err)
}
if len(*asked) != 0 || len(pendingOf(t, open)) != 0 {
t.Fatalf("an unknown module asked %v, pending %+v", *asked, pendingOf(t, open))
}
}
// Several modules in one act, one of them not registered: nothing is assigned and nothing kept pending, and
// each missing one is said.
func TestAnActWithAModuleNotRegisteredIsRefusedWhole(t *testing.T) {
open := aCatalogueMesh(t)
ctx := t.Context()
asksWithPaths(t)
register(t, open, catalogue.Manifest{Module: "known", Version: "1"})
mergeAdding(t, open, "modules/sensors", "3da80a4b00aa")
_, err := assign(ctx, open, "laptop", "known", "sensors")
if err == nil || !strings.Contains(err.Error(), "nothing was assigned") || !strings.Contains(err.Error(), "being built") {
t.Fatalf("an act naming a module in flight: %v", err)
}
if slices.Contains(assignedTo(t, open, "laptop"), "known") || len(pendingOf(t, open)) != 0 {
t.Fatal("an act refused whole was made in part")
}
}
// unassign withdraws a waiting pending assignment; the build goes on and registers the module assigned
// nowhere: a withdrawn row is never claimed.
func TestUnassignWithdrawsAPendingAssignment(t *testing.T) {
open := aCatalogueMesh(t)
ctx := t.Context()
asksWithPaths(t)
mergeAdding(t, open, "modules/sensors", "3da80a4b00aa")
if _, err := assign(ctx, open, "laptop", "sensors"); err != nil {
t.Fatal(err)
}
said, err := unassign(ctx, open, "laptop", "sensors")
if err != nil || !strings.Contains(said, "withdrawn") {
t.Fatalf("unassign of a pending assignment: %q %v", said, err)
}
if err := (builds{open.inventory, open}).Built(ctx, outcome("b-modules/sensors", "modules/sensors",
"3da80a4b00aa", "sensors", "")); err != nil {
t.Fatal(err)
}
settlePending(ctx, open, time.Now())
if slices.Contains(assignedTo(t, open, "laptop"), "sensors") {
t.Fatal("a withdrawn assignment was made")
}
}
// **A claim is never overridden** (review point 1): a pending assignment being made is claimed
// (waiting → applying); an unassign meanwhile is refused and says so, and leaves the claim; making a row
// already withdrawn makes nothing. A claim left by a controller that stopped is settled by the tick from what
// the mesh holds: back to waiting when the module is not assigned, applied when it is.
func TestAWithdrawalNeverOverridesAClaim(t *testing.T) {
open := aCatalogueMesh(t)
ctx := t.Context()
inv := open.inventory
asksWithPaths(t)
mergeAdding(t, open, "modules/sensors", "3da80a4b00aa")
if _, err := assign(ctx, open, "laptop", "sensors"); err != nil {
t.Fatal(err)
}
p := pendingOf(t, open)[0]
if ok, err := inv.ClaimPending(ctx, p.ID); err != nil || !ok {
t.Fatalf("claim: %v %v", ok, err)
}
_, err := unassign(ctx, open, "laptop", "sensors")
if err == nil || !strings.Contains(err.Error(), "is being assigned sensors now") {
t.Fatalf("unassign of a claimed pending assignment: %v", err)
}
if got := pendingOf(t, open)[0]; got.State != inventory.PendingApplying {
t.Fatalf("the claim was overridden: %+v", got)
}
// Pending(zero) is the waiting ones alone (review point 7).
if waiting, err := inv.Pending(ctx, time.Time{}); err != nil || len(waiting) != 0 {
t.Fatalf("Pending(zero) read a claimed row: %+v %v", waiting, err)
}
// Left by a controller that stopped: back to waiting, since sensors is not assigned.
settlePending(ctx, open, time.Now().Add(claimStaleAfter+time.Minute))
if got := pendingOf(t, open)[0]; got.State != inventory.PendingWaiting {
t.Fatalf("a stale claim was not released: %+v", got)
}
// Withdrawn, then the build registers it: applyPending finds nothing to claim.
if _, err := unassign(ctx, open, "laptop", "sensors"); err != nil {
t.Fatal(err)
}
register(t, open, catalogue.Manifest{Module: "sensors", Version: "1"})
if line := applyPending(ctx, open, p, "b-modules/sensors"); line != "" {
t.Fatalf("a withdrawn pending assignment was made: %s", line)
}
if slices.Contains(assignedTo(t, open, "laptop"), "sensors") {
t.Fatal("a withdrawn pending assignment was assigned")
}
// A stale claim whose module was assigned is applied.
q, err := inv.RecordPending(ctx, inventory.PendingAssignment{Node: "anchor", Module: "sensors", Build: "b-x"})
if err != nil {
t.Fatal(err)
}
if ok, _ := inv.ClaimPending(ctx, q.ID); !ok {
t.Fatal("claim")
}
if _, err := inv.Assign(ctx, "anchor", "sensors"); err != nil {
t.Fatal(err)
}
settlePending(ctx, open, time.Now().Add(claimStaleAfter+time.Minute))
rows, _ := inv.PendingFor(ctx, "anchor", "sensors")
if len(rows) != 1 || rows[0].State != inventory.PendingApplied {
t.Fatalf("a stale claim of an assigned module: %+v", rows)
}
}
// A build that says nothing within the bound: the pending assignment expires at the next tick, saying so,
// rather than waiting for ever; and a module registered by a process that kept no pending assignment is
// assigned at the next tick.
func TestAPendingAssignmentNeverWaitsSilently(t *testing.T) {
open := aCatalogueMesh(t)
ctx := t.Context()
inv := open.inventory
for _, r := range []inventory.BuildRequest{
{ID: "build-lost", Repository: "novox/mesh-catalog", Seat: "git", Path: "modules/lost", Ref: "main",
For: "merge", At: time.Now().Add(-buildRequestBound / 2)},
{ID: "build-heard", Repository: "novox/mesh-catalog", Seat: "git", Path: "modules/heard", Ref: "main",
For: "merge", At: time.Now()},
} {
if err := inv.RecordBuildRequest(ctx, r); err != nil {
t.Fatal(err)
}
}
for _, m := range []string{"lost", "heard"} {
if _, err := assign(ctx, open, "laptop", m); err != nil {
t.Fatal(err)
}
}
register(t, open, catalogue.Manifest{Module: "heard", Version: "1"})
settlePending(ctx, open, time.Now().Add(buildRequestBound))
if !slices.Contains(assignedTo(t, open, "laptop"), "heard") {
t.Fatal("a module registered meanwhile was not assigned at the next tick")
}
byModule := map[string]inventory.PendingAssignment{}
for _, p := range pendingOf(t, open) {
byModule[p.Module] = p
}
if p := byModule["lost"]; p.State != inventory.PendingExpired || !strings.Contains(p.Note, "no outcome of build build-lost") {
t.Fatalf("a build that said nothing: %+v", p)
}
if p := byModule["heard"]; p.State != inventory.PendingApplied {
t.Fatalf("a module registered meanwhile: %+v", p)
}
// Ended rows are deleted after KeptFor (review point 7); open ones never, nor one whose condition is
// still open (second review, bug B).
if _, err := inv.RecordPending(ctx, inventory.PendingAssignment{Node: "anchor", Module: "lost", Build: "build-lost"}); err != nil {
t.Fatal(err)
}
said := settlePending(ctx, open, time.Now().Add(inventory.KeptFor+time.Hour))
rows, err := inv.Pending(ctx, time.Unix(0, 0))
if err != nil {
t.Fatal(err)
}
for _, r := range rows {
if !r.Open() && (r.Raised == nil || r.Cleared != nil) {
t.Fatalf("an ended pending assignment with no open condition outlived %s: %+v (%v)", inventory.KeptFor, r, said)
}
}
if !strings.Contains(strings.Join(said, "\n"), "deleted 1 pending assignment(s)") {
t.Fatalf("the applied row was not deleted: %v", said)
}
}
// The seat's assign passes build on; unassign takes none.
func TestTheSeatsAssignPassesBuild(t *testing.T) {
argv, err := argvFor("assign", map[string]any{"node": "g14", "module": "sensors", "build": "true"})
if err != nil || !slices.Equal(argv, []string{"assign", "g14", "sensors", "--build"}) {
t.Fatalf("assign with build: %v %v", argv, err)
}
if _, err := argvFor("unassign", map[string]any{"node": "g14", "module": "sensors", "build": "true"}); err == nil {
t.Fatal("unassign took build")
}
}
// The condition's words are plain.
func TestAnAssignmentNotMadeIsSaidPlainly(t *testing.T) {
o := pendingObservation(inventory.PendingAssignment{Node: "g14", Module: "sensors",
Note: "the pending assignment of sensors to g14 expired: build b-1 of modules/sensors failed: boom"})
w := plainWordings[kindPendingEnded](o)
if why, ok := conditions.PlainWords(w, "g14"); !ok {
t.Fatalf("not plain: %s %+v", why, w)
}
if w.Headline != "sensors was not put on g14" {
t.Fatalf("headline %q", w.Headline)
}
}
// **One row's condition is never another's** (second review, bug A). The sequence: the first pending
// assignment's build fails and the tick raises it; the person assigns again with build "true", a second
// pending assignment; its build fails too before the next tick. That tick raises the second, and must not
// clear it by judging the first answered by the second: each row has its own condition, and a newer row that
// itself ended unmade answers nothing.
func TestASecondAssignmentNotMadeKeepsItsCondition(t *testing.T) {
open := aCatalogueMesh(t)
ctx := t.Context()
asked := 0
was := askABuild
askABuild = func(context.Context, buildSource, string, string) (string, error) {
asked++
return fmt.Sprintf("b-%d", asked), nil
}
t.Cleanup(func() { askABuild = was })
mergeAdding(t, open, "modules/sensors", "3da80a4b00aa")
if _, err := assign(ctx, open, "laptop", "sensors"); err != nil {
t.Fatal(err)
}
failedBuild(t, open, "b-1")
settlePending(ctx, open, time.Now())
first := pendingOf(t, open)[0]
if !conditionOpen(t, first) {
t.Fatal("the first assignment not made raised nothing")
}
if _, err := assignWith(ctx, open, "laptop", assignOptions{Build: true}, "sensors"); err != nil {
t.Fatal(err)
}
failedBuild(t, open, "b-2")
settlePending(ctx, open, time.Now())
var second inventory.PendingAssignment
for _, p := range pendingOf(t, open) {
if p.Build == "b-2" {
second = p
}
}
if second.State != inventory.PendingExpired {
t.Fatalf("the second: %+v", second)
}
if !conditionOpen(t, second) {
t.Fatal("the second assignment's condition was cleared in the tick that raised it")
}
if !conditionOpen(t, rowOf(t, open, first.ID)) {
t.Fatal("the first was judged answered by a second that itself was not made")
}
// unassign takes both back, and the next tick clears both.
if said, err := unassign(ctx, open, "laptop", "sensors"); err != nil || !strings.Contains(said, "b-1") ||
!strings.Contains(said, "b-2") {
t.Fatalf("unassign: %q %v", said, err)
}
settlePending(ctx, open, time.Now())
if conditionOpen(t, first) || conditionOpen(t, second) {
t.Fatal("a condition stayed open after both were taken back")
}
}
// **A raised row outlives the pruning until its condition clears** (second review, bug B); and a machine's
// removal, which takes its rows with it, clears their conditions at the next tick.
func TestARaisedRowIsKeptUntilItsConditionClears(t *testing.T) {
open := aCatalogueMesh(t)
ctx := t.Context()
inv := open.inventory
asksWithPaths(t)
mergeAdding(t, open, "modules/sensors", "3da80a4b00aa")
if _, err := assign(ctx, open, "laptop", "sensors"); err != nil {
t.Fatal(err)
}
failedBuild(t, open, "b-modules/sensors")
settlePending(ctx, open, time.Now())
row := pendingOf(t, open)[0]
settlePending(ctx, open, time.Now().Add(inventory.KeptFor+time.Hour))
if rowOf(t, open, row.ID).Raised == nil || !conditionOpen(t, row) {
t.Fatal("a raised row was pruned, or its condition closed, while the condition was open")
}
if _, err := unassign(ctx, open, "laptop", "sensors"); err != nil {
t.Fatal(err)
}
settlePending(ctx, open, time.Now())
settlePending(ctx, open, time.Now().Add(inventory.KeptFor+time.Hour))
if known, err := inv.PendingKnown(ctx, []int64{row.ID}); err != nil || known[row.ID] {
t.Fatalf("a cleared row outlived %s: %v", inventory.KeptFor, err)
}
// On anchor, then anchor removed: the row goes, and its condition with it at the next tick.
if _, err := assign(ctx, open, "anchor", "sensors"); err == nil {
t.Fatal("known and not built was not refused")
}
q, err := inv.RecordPending(ctx, inventory.PendingAssignment{Node: "anchor", Module: "sensors", Build: "b-x"})
if err != nil {
t.Fatal(err)
}
if _, err := inv.SettlePending(ctx, q.ID, inventory.PendingWaiting, inventory.PendingExpired, "boom"); err != nil {
t.Fatal(err)
}
settlePending(ctx, open, time.Now())
q.Node = "anchor"
if !conditionOpen(t, q) {
t.Fatal("not raised")
}
if _, err := inv.RemoveNodeForTest(ctx, "anchor"); err != nil {
t.Fatal(err)
}
settlePending(ctx, open, time.Now())
if conditionOpen(t, q) {
t.Fatal("a removed machine's condition stayed open")
}
}
// **A claim left by a controller that stopped** (second review): the tick settles it from what the mesh
// holds — back to waiting when the module is not assigned there, applied when it is — and leaves a fresh
// claim alone.
func TestAStaleClaimIsSettledFromWhatTheMeshHolds(t *testing.T) {
open := aCatalogueMesh(t)
ctx := t.Context()
inv := open.inventory
claimed := func(node string) inventory.PendingAssignment {
p, err := inv.RecordPending(ctx, inventory.PendingAssignment{Node: node, Module: "sensors", Build: "b-x"})
if err != nil {
t.Fatal(err)
}
if ok, err := inv.ClaimPending(ctx, p.ID); err != nil || !ok {
t.Fatalf("claim: %v %v", ok, err)
}
return p
}
notAssigned, assigned := claimed("laptop"), claimed("anchor")
register(t, open, catalogue.Manifest{Module: "unrelated", Version: "1"})
if err := inv.RegisterModule(ctx, catalogue.Manifest{Module: "sensors", Version: "1"}, inventory.Source{}); err != nil {
t.Fatal(err)
}
if _, err := inv.Assign(ctx, "anchor", "sensors"); err != nil {
t.Fatal(err)
}
settlePending(ctx, open, time.Now())
if rowOf(t, open, notAssigned.ID).State != inventory.PendingApplying || rowOf(t, open, assigned.ID).State != inventory.PendingApplying {
t.Fatal("a fresh claim was settled")
}
settlePending(ctx, open, time.Now().Add(claimStaleAfter+time.Minute))
if got := rowOf(t, open, assigned.ID); got.State != inventory.PendingApplied {
t.Fatalf("a stale claim of an assigned module: %+v", got)
}
// Released to waiting, and in the same tick made, since sensors is registered now.
if got := rowOf(t, open, notAssigned.ID); got.State != inventory.PendingApplied ||
!slices.Contains(assignedTo(t, open, "laptop"), "sensors") {
t.Fatalf("a stale claim of a module not assigned: %+v", got)
}
}
// **A waited build whose asker stopped waiting is asked, outcome unknown** (second review): it may still run,
// so it reads as in flight until its outcome or its bound, never as not asked; a build never handed over
// is not asked.
func TestABuildNoLongerWaitedForIsStillInFlight(t *testing.T) {
open := aCatalogueMesh(t)
ctx := t.Context()
inv := open.inventory
for _, id := range []string{"build-timeout", "build-nothandedover"} {
dir := "modules/" + strings.TrimPrefix(id, "build-")
if err := inv.RecordBuildRequest(ctx, inventory.BuildRequest{ID: id, Repository: "novox/mesh-catalog",
Seat: "git", Path: dir, Ref: "main", For: "build"}); err != nil {
t.Fatal(err)
}
}
markWaitFailed(ctx, "build-timeout", errors.New("no build machine answered within 10m0s"))
markWaitFailed(ctx, "build-nothandedover", fmt.Errorf("%w: cannot submit a build: nats: timeout", link.ErrNotHandedOver))
said, err := assign(ctx, open, "laptop", "timeout")
if err != nil || !strings.Contains(said, "being built") || !strings.Contains(said, "kept as pending") {
t.Fatalf("a build no longer waited for: %q %v", said, err)
}
if _, err := assign(ctx, open, "laptop", "nothandedover"); err == nil || !strings.Contains(err.Error(), "was not handed over") {
t.Fatalf("a build never handed over: %v", err)
}
settlePending(ctx, open, time.Now().Add(buildRequestBound+time.Minute))
rows, _ := inv.PendingFor(ctx, "laptop", "timeout")
if len(rows) != 1 || rows[0].State != inventory.PendingExpired || !strings.Contains(rows[0].Note, "its asker stopped waiting") {
t.Fatalf("past the bound: %+v", rows)
}
}
+13
View File
@@ -153,6 +153,19 @@ var plainWordings = map[string]func(conditions.Observation) words{
"asleep is normal.", m),
Resolved: m + " can be reached again"}
}),
// A pending assignment that was not made (novox/hq issue 325, ADR 0261).
kindPendingEnded: worded(func(o conditions.Observation) words {
m := machineOr(o, "a machine")
module := idPart(o, 1)
if module == "" {
module = "a module"
}
return words{Headline: module + " was not put on " + m,
Needs: "decide whether to build it again or take the assignment back; the details say why.",
Explanation: fmt.Sprintf("An assignment of %s to %s waited for its build, and the build did not "+
"register it, so it was not made.", module, m),
Resolved: "Resolved: " + module + " on " + m + " is settled"}
}),
kindBindingKept: worded(func(o conditions.Observation) words {
m := machineOr(o, "a machine")
return words{Headline: "A module's data source is held on " + m,
+4
View File
@@ -165,6 +165,10 @@ func serve(ctx context.Context) (err error) {
// Open plans move on a timer as well as on outcomes (novox/hq ADR 0162): a tier waiting for
// machines to report moves when they have, and a plan left by a replaced controller resumes.
go planTicker(ctx, open)
// And the pending assignments settled on a tick of their own (novox/hq ADR 0261): made once their module
// is registered, ended with why when its build will not register it, raised and cleared as conditions.
// Never by a read.
go settlingPending(ctx, open)
// The durations the core's bounds are set from are kept a month (novox/hq to-be 45 Phase 0).
go forgettingOldDurations(ctx, inv)
// And what the catalogue decided a build meant. The builder's own result is already handled
+7 -1
View File
@@ -406,6 +406,8 @@ func rebuildCommand(ctx context.Context, args []string) error {
if err != nil {
return err
}
recordAsked(ctx, inventory.BuildRequest{ID: id, Repository: source.Repository, Seat: source.Seat, Path: path,
Ref: ref, For: "rebuild"})
fmt.Printf("rebuild asked as %s\n", id)
return nil
}
@@ -490,7 +492,11 @@ func replayCommand(ctx context.Context, args []string) error {
if err := replayRefusal(b, module, history, *register, *older, outstanding); err != nil {
return err
}
id, err := buildOneAsked(ctx, source, b.Path, b.Commit, 0, !*register)
asker := ""
if *register {
asker = "replay"
}
id, err := buildOneAsked(ctx, source, b.Path, b.Commit, 0, !*register, asker)
if err != nil {
return err
}
+23
View File
@@ -101,6 +101,24 @@ type meshStatus struct {
// Overflowing is every module whose identity overflows the bound of a provision it requires, and
// so is left out of its provider's grants (novox/hq ADR 0225). Absent when every identity fits.
Overflowing []catalogue.Overflow `json:"overflowing,omitempty"`
// Pending is every assignment waiting for its module's build to register it, and every one ended in
// the last day with what ended it (novox/hq issue 325). Absent when there is none; PendingUnread says
// why they could not be read.
Pending []pendingStatus `json:"pending,omitempty"`
PendingUnread string `json:"pendingUnread,omitempty"`
}
// pendingStatus is one pending assignment as status says it.
type pendingStatus struct {
Node string `json:"node"`
Module string `json:"module"`
State string `json:"state"`
Build string `json:"build"`
Repository string `json:"repository,omitempty"`
Path string `json:"path,omitempty"`
Since time.Time `json:"since"`
Settled *time.Time `json:"settled,omitempty"`
Note string `json:"note,omitempty"`
}
// machineFiltered is one rule set on a converged machine that the mesh did not write and that
@@ -240,6 +258,11 @@ func statusAsJSON(asked answers) ([]byte, error) {
out.Conditions, out.ConditionsUnread = inBrief(asked.conditions), asked.conditionsUnread
out.Failing = providerStandings(asked.conditions)
out.Overflowing = asked.overflowing
out.PendingUnread = asked.pendingUnread
for _, p := range asked.pending {
out.Pending = append(out.Pending, pendingStatus{Node: p.Node, Module: p.Module, State: p.State,
Build: p.Build, Repository: p.Repository, Path: p.Path, Since: p.Since, Settled: p.Settled, Note: p.Note})
}
for name := range asked.refused {
out.Unresolved = append(out.Unresolved, machineUnresolved{
Node: name, Problem: asked.refused[name]})
+12 -2
View File
@@ -337,8 +337,16 @@ func askTier(ctx context.Context, inv *inventory.Inventory, p *inventory.Plan) e
// askABuild is how a plan asks for one build, not waited for, and learns the id it asked under. A
// variable so a test of what a plan does around an ask needs no build machine.
var askABuild = func(ctx context.Context, source buildSource, path, ref string) (string, error) {
return buildOneAsked(ctx, source, path, ref, 0, false)
//
// Set in init, not where it is declared: an assignment pending on a build asks for one too (novox/hq issue
// 325), and a build's outcome makes pending assignments, so the two refer to each other.
var askABuild func(ctx context.Context, source buildSource, path, ref string) (string, error)
func init() {
askABuild = func(ctx context.Context, source buildSource, path, ref string) (string, error) {
// Kept by each asker with its own name once the ask is made (novox/hq issue 325).
return buildOneAsked(ctx, source, path, ref, 0, false, "")
}
}
// askModule asks the build machine for one module of a plan and marks it asked, with the id it was
@@ -369,6 +377,8 @@ func askModule(ctx context.Context, p *inventory.Plan, name string, byName map[s
p.Note = fmt.Sprintf("%s could not be asked for: %v", name, err)
return
}
recordAsked(ctx, inventory.BuildRequest{ID: id, Repository: e.Source.Repository, Seat: e.Source.Seat,
Path: e.Source.Path, Ref: followedBranch(e.Source.Ref), Commit: p.Commit, For: "plan"})
state.State = "asked"
state.AskedAt = &now
state.Build = id
+60
View File
@@ -0,0 +1,60 @@
package main
import (
"encoding/json"
"slices"
"strings"
"testing"
"github.com/novox/mesh-controller/internal/catalogue"
"github.com/novox/mesh-controller/internal/inventory"
"github.com/novox/mesh-controller/internal/link"
)
// novox/hq issue 325, replayed with only what the controller had before its fix, so it can be laid over the
// older commit. On 2026-10-08 the catalogue's pull request adding `sensors` merged, and the merge asked for
// its build (issue 300). About a minute later `assign g14 sensors` answered "no module of that name:
// sensors"; a few minutes later the same call worked, because the build had registered the module. The
// answer read as "nobody registered it". An assignment made while the build runs says the build — where it
// is from, since when, its id — is kept, and is made when the build registers the module.
func TestReplay325(t *testing.T) {
open := aMesh(t)
ctx := t.Context()
if err := open.inventory.RegisterModule(ctx, catalogue.Manifest{Module: "networkmanager", Version: "1"},
inventory.Source{Repository: "novox/mesh-catalog", Seat: "git", Path: "modules/networkmanager", Ref: "main",
BuiltFrom: "c0", Head: "c0"}); err != nil {
t.Fatal(err)
}
asksWithPaths(t)
m := link.SourceMoved{Owner: "novox", Repo: "mesh-catalog", Base: "main", Commit: "5e450125aa",
Paths: []string{"modules/sensors/module.json", "modules/sensors/cmd/sensors/main.go"},
ModuleDirs: []string{"modules/sensors"}, ModuleDirsSaid: true}
if err := (following{open: open}).SourceMoved(ctx, m); err != nil {
t.Fatal(err)
}
said, err := assign(ctx, open, "laptop", "sensors")
if err != nil {
t.Fatalf("assign while the merge's build runs was refused: %v", err)
}
for _, want := range []string{"being built", "novox/mesh-catalog", "modules/sensors", "b-modules/sensors"} {
if !strings.Contains(said, want) {
t.Fatalf("the answer does not say %q:\n%s", want, said)
}
}
manifest, _ := json.Marshal(catalogue.Manifest{Module: "sensors", Version: "1"})
if err := (builds{open.inventory, open}).Built(ctx, link.BuildResult{ID: "b-modules/sensors",
Repository: "novox/mesh-catalog", Path: "modules/sensors", Ref: "main", Commit: "5e450125aa", On: "anchor",
Module: "sensors", Manifest: manifest,
Source: &link.SourceOnSeat{Repository: "novox/mesh-catalog", Seat: "git"}}); err != nil {
t.Fatal(err)
}
assigned, err := open.inventory.Assigned(ctx, "laptop")
if err != nil {
t.Fatal(err)
}
if !slices.Contains(assigned, "sensors") {
t.Fatalf("the build registered sensors and laptop runs %v", assigned)
}
}
+7 -1
View File
@@ -418,7 +418,13 @@ func (a *verbArguments) commandLine() ([]string, error) {
}
// Several modules comma-separated, judged as one act (novox/hq ADR 0207): the holders of
// the seats that apply resources depend on each other and go on together.
return append([]string{verb, str("node")}, splitModules(str("module"))...), nil
argv := append([]string{verb, str("node")}, splitModules(str("module"))...)
// A module known and not built: its build asked for, the assignment pending on it (novox/hq
// issue 325). unassign declares no build, so a call giving it is refused before this.
if verb == "assign" && on("build") {
argv = append(argv, "--build")
}
return argv, nil
case "pin":
if err := need("node", "provision", "from", "module"); err != nil {
return nil, err
+56 -1
View File
@@ -143,6 +143,8 @@ func printStatus(asked answers) error {
len(quiet), strings.Join(said, "\n "))
}
printPending(asked.pending, asked.pendingUnread)
if open, late := openPlans(asked.plans, time.Now(), asked.paused, asked.tierBounds); len(open) > 0 {
fmt.Printf("%d plan(s) open", len(open))
if late > 0 {
@@ -466,6 +468,12 @@ func theThreeQuestions(ctx context.Context, open *stores) (answers, error) {
if out.conditions, err = openConditions(ctx); err != nil {
out.conditionsUnread = err.Error()
}
// And the assignments waiting for their module's build (novox/hq issue 325, ADR 0261), read only: the
// controller's tick settles them, never a read.
if out.pending, err = inv.Pending(ctx, time.Now().Add(-pendingShownFor)); err != nil {
out.pendingUnread = err.Error()
err = nil
}
out.plans, err = inv.RecentPlans(ctx, 5)
if err != nil {
return answers{}, err
@@ -598,7 +606,54 @@ func (a answers) well() bool {
return len(a.wrong) == 0 && len(a.quiet) == 0 && len(a.behind) == 0 &&
len(a.waiting) == 0 && len(a.refused) == 0 && a.network == "" && len(a.untaken) == 0 &&
len(a.filtered) == 0 && len(a.unheld) == 0 && len(a.overflowing) == 0 &&
len(a.conditions) == 0 && a.conditionsUnread == ""
len(a.conditions) == 0 && a.conditionsUnread == "" && a.pendingUnread == "" && !pendingOpen(a.pending)
}
// pendingOpen is whether any pending assignment is still to be made. One that ended without being made is
// not counted here: it is a condition, open until it is answered (ADR 0261).
func pendingOpen(pending []inventory.PendingAssignment) bool {
for _, p := range pending {
if p.Open() {
return true
}
}
return false
}
// printPending says the assignments waiting for their module's build, and those ended in the last day
// (novox/hq issue 325).
func printPending(pending []inventory.PendingAssignment, unread string) {
if unread != "" {
fmt.Printf("the assignments waiting for a build could not be read: %s\n\n", unread)
return
}
var waiting, ended []inventory.PendingAssignment
for _, p := range pending {
if p.Open() {
waiting = append(waiting, p)
} else {
ended = append(ended, p)
}
}
if len(waiting) > 0 {
fmt.Printf("%d assignment(s) wait for their module's build to register it:\n", len(waiting))
for _, p := range waiting {
fmt.Printf(" %-12s %-20s build %s of %s %s, since %s\n", p.Node, p.Module, p.Build, p.Repository,
orRoot(p.Path), clock(p.Since))
}
fmt.Printf("\n each is made when its build registers the module; `unassign <node> <module>` withdraws it\n\n")
}
if len(ended) > 0 {
fmt.Printf("%d pending assignment(s) ended in the last day:\n", len(ended))
for _, p := range ended {
at := ""
if p.Settled != nil {
at = clock(*p.Settled)
}
fmt.Printf(" %-12s %-20s %-9s %s\n %-12s %s\n", p.Node, p.Module, p.State, at, "", p.Note)
}
fmt.Println()
}
}
// hostSplit is which machines report which host version, for every version more than one machine
+11 -2
View File
@@ -360,7 +360,7 @@ func (f following) SourceMoved(ctx context.Context, m link.SourceMoved) error {
}
// **A new module is built and registered, and sent nowhere** (novox/hq issue 300): the delivery plan
// says so, and assigning it is a person's act, which needs it registered first.
built := askNewModules(ctx, m, from, added)
built := askNewModules(ctx, inv, m, from, added)
moved := append(append([]inventory.Entry{}, touched...), packaging...)
if len(moved) == 0 {
if len(built) > 0 {
@@ -491,7 +491,8 @@ func (f following) SourceMoved(ctx context.Context, m link.SourceMoved) error {
// new module either. Outside the plan: the plan walks modules the catalogue holds, and a new one has no
// machine to send to and nothing standing on it. What could not be asked is said, and left to `build`.
// Returns the directories asked for.
func askNewModules(ctx context.Context, m link.SourceMoved, from []inventory.Entry, added []string) []string {
func askNewModules(ctx context.Context, inv *inventory.Inventory, m link.SourceMoved, from []inventory.Entry,
added []string) []string {
if len(added) == 0 || len(from) == 0 {
return nil
}
@@ -503,11 +504,19 @@ func askNewModules(ctx context.Context, m link.SourceMoved, from []inventory.Ent
path = ""
}
id, err := askABuild(ctx, source, path, m.Base)
// Kept, asked or not, so `assign` can tell a build in flight, or a merge that could not ask for one,
// from a module nobody ever heard of (novox/hq issue 325).
request := inventory.BuildRequest{ID: id, Repository: source.Repository, Seat: source.Seat, Path: path,
Ref: m.Base, Commit: m.Commit, For: "merge"}
if err != nil {
request.ID = fmt.Sprintf("not-asked-%s-%s", short(m.Commit), strings.ReplaceAll(dir, "/", "-"))
request.NotAsked = err.Error()
recordBuildRequest(ctx, inv, request)
fmt.Printf(" %s is a new module in %s/%s and could not be built: %v — a hand `build` of it "+
"asks again\n", dir, m.Owner, m.Repo, err)
continue
}
recordBuildRequest(ctx, inv, request)
fmt.Printf(" %s is a new module in %s/%s: build %s asked at %s; registered when it lands, assigned "+
"nowhere\n", dir, m.Owner, m.Repo, id, m.Base)
asked = append(asked, dir)
+10 -2
View File
@@ -178,9 +178,17 @@ var ControllerVerbs = []Verb{
"diff": "\"true\": what a push would change on this machine, against what it was last sent; not with files",
}, []string{"node"}, "files", "diff")},
{Name: "assign", Description: "Put a module on a machine. Refused with the mesh's own words when it cannot resolve there, " +
"or when a seat its resources are applied through is held by nothing on the machine (novox/hq ADR 0207).",
"or when a seat its resources are applied through is held by nothing on the machine (novox/hq ADR 0207). " +
"A module not registered yet is looked for (novox/hq issue 325, ADR 0261): while its build runs — a merge " +
"that adds a module asks for one — the assignment is kept as pending and the controller makes it when the " +
"build registers it; `status` lists it, and if the build fails it expires with why and raises a condition. " +
"A module known and not built is said, and build asks for its build; a name with no module and no build " +
"request kept is refused with the closest names.",
Input: schema(map[string]string{"node": "the machine's name",
"module": "the module's name; several comma-separated are judged together"}, []string{"node", "module"})},
"module": "the module's name; several comma-separated are judged together",
"build": "\"true\": for a module known and not built (its last build failed, or was never asked), ask " +
"for its build and keep this assignment pending until the build registers it"},
[]string{"node", "module"}, "build")},
{Name: "unassign", Description: "Take a module off a machine. Refused when it holds a seat a module left there depends on.",
Input: schema(map[string]string{"node": "the machine's name",
"module": "the module's name; several comma-separated are judged together"}, []string{"node", "module"})},
+10
View File
@@ -68,3 +68,13 @@ func ForTest(t *testing.T) *Inventory {
}
return inv
}
// RemoveNodeForTest removes a machine's record as a person removing it from the store would: no verb of the
// mesh removes one yet, and what hangs off it must go with it.
func (i *Inventory) RemoveNodeForTest(ctx context.Context, name string) (bool, error) {
tag, err := i.store.Pool().Exec(ctx, `delete from node where name = $1`, name)
if err != nil {
return false, err
}
return tag.RowsAffected() == 1, nil
}
@@ -0,0 +1,52 @@
-- An assignment waits for its module's build (novox/hq issue 325, ADR 0261).
--
-- A merge that adds a module asks for its build (issue 300), and the module is registered when the build's
-- outcome is taken in, minutes later. An `assign` in between was answered "no module of that name", which
-- read as "nobody registered it". The controller now keeps every build it asks for, so `assign` can tell a
-- build in flight from a module it never heard of, and an assignment made while the build runs is kept as
-- pending and made by the controller when the build registers the module, or ended with the reason.
-- Every build the controller asked for: a merge's new module, a plan's tier, a person's `build`, an `assign`
-- with build. Kept once the ask is made, never before. Its outcome is the `build` row of the same id;
-- an ask without one is still running, or was lost. not_asked is why the ask could not be handed over; for a
-- merge that could not ask, the id is the controller's and no build's. outcome_unknown is why an asker that
-- handed it over stopped waiting: the build may still run, and is read as in flight until its outcome or
-- its bound.
create table build_request (
id text primary key,
repository text not null,
seat text not null default '',
source_path text not null default '',
ref text not null default '',
commit_hash text not null default '',
asked_for text not null default '',
not_asked text,
outcome_unknown text,
asked_at timestamptz not null default now()
);
create index build_request_asked_at on build_request (asked_at);
-- An assignment kept until its module is registered (ADR 0261). waiting until the build registers the
-- module; applying while the controller makes it, under the machine's hold; then applied, refused (the
-- assignment was refused then), expired (the build failed, was not registered, or said nothing within its
-- bound) or withdrawn (`unassign`). raised_at is when an expiry or refusal was raised as a condition,
-- acknowledged_at when a person took it back with `unassign`, cleared_at when its condition was cleared.
-- Gone with its machine. Ended rows are deleted after 30 days, once any condition raised for them is cleared.
create table pending_assignment (
id bigserial primary key,
node uuid not null references node (id) on delete cascade,
module text not null,
build_id text not null,
repository text not null default '',
source_path text not null default '',
since timestamptz not null default now(),
state text not null default 'waiting'
check (state in ('waiting', 'applying', 'applied', 'refused', 'expired', 'withdrawn')),
note text not null default '',
settled_at timestamptz,
raised_at timestamptz,
acknowledged_at timestamptz,
cleared_at timestamptz
);
create unique index pending_assignment_open on pending_assignment (node, module)
where state in ('waiting', 'applying');
+387
View File
@@ -0,0 +1,387 @@
package inventory
import (
"context"
"errors"
"path"
"strings"
"time"
"github.com/jackc/pgx/v5"
)
// What the controller asked to be built, and the assignments waiting for it (novox/hq issue 325, ADR 0261).
//
// A build's outcome was kept, and its asking was not: between a merge asking for a new module and the
// outcome registering it, the mesh had no record that anything was coming, and `assign` answered "no module
// of that name" for a module minutes from existing. Kept here, a build request lets `assign` say which case
// it is in, and an assignment made while the build runs is kept until the build registers the module.
// KeptFor is how long a build request, and an ended pending assignment, is kept.
const KeptFor = 30 * 24 * time.Hour
// BuildRequest is one build the controller asked for.
type BuildRequest struct {
// ID is the build's id, as its outcome will carry it; for a merge whose ask could not be made, an id
// of the controller's making that no build carries.
ID string
// Repository and Seat are the source as the mesh spells it (novox/hq ADR 0111): a path on the seat's
// holder, or a URL with no seat. Path is the module's directory in it, Ref the branch or commit asked.
Repository, Seat, Path, Ref string
// Commit is the commit that made the ask, where one did: a merge's.
Commit string
// For says who asked: "merge", "plan", "build" (a person), "assign".
For string
// NotAsked is why the ask could not be handed over; empty when it was.
NotAsked string
// OutcomeUnknown is why an asker that handed the build over stopped waiting for it: asked, outcome
// unknown. Read as in flight until its outcome or its bound.
OutcomeUnknown string
At time.Time
}
// Name is the module this request is expected to register, read from its directory: the last element of
// its path, or of its repository for a module at the repository's root. A guess until the build says, and
// the one a person reads too: `modules/sensors` is sensors.
func (a BuildRequest) Name() string {
if p := strings.Trim(a.Path, "/"); p != "" && p != "." {
return path.Base(p)
}
return strings.TrimSuffix(path.Base(strings.TrimRight(a.Repository, "/")), ".git")
}
// RequestOutcome is a build request beside what came of it, where anything did.
type RequestOutcome struct {
BuildRequest
// Heard is whether an outcome of this id was taken in; Failed is its builder's words when it failed,
// Module what the build said it built, HeardAt when it was taken in.
Heard bool
Failed string
Module string
HeardAt time.Time
}
// RecordBuildRequest keeps one build request, once it is asked. Idempotent on the id. Requests older than
// KeptFor are deleted in the same act.
func (i *Inventory) RecordBuildRequest(ctx context.Context, a BuildRequest) error {
at := a.At
if at.IsZero() {
at = time.Now().UTC()
}
var notAsked *string
if a.NotAsked != "" {
notAsked = &a.NotAsked
}
if _, err := i.store.Pool().Exec(ctx,
`insert into build_request (id, repository, seat, source_path, ref, commit_hash, asked_for, not_asked, asked_at)
values ($1, $2, $3, $4, $5, $6, $7, $8, $9)
on conflict (id) do nothing`,
a.ID, a.Repository, a.Seat, strings.Trim(a.Path, "/"), a.Ref, a.Commit, a.For, notAsked, at); err != nil {
return err
}
_, err := i.store.Pool().Exec(ctx, `delete from build_request where asked_at < $1`, time.Now().Add(-KeptFor))
return err
}
// MarkNotAsked says a kept build request was never handed over: the words are kept, unless its outcome was
// heard first.
func (i *Inventory) MarkNotAsked(ctx context.Context, id, why string) error {
_, err := i.store.Pool().Exec(ctx,
`update build_request set not_asked = $2
where id = $1 and not exists (select 1 from build where build.id = $1)`, id, why)
return err
}
// MarkOutcomeUnknown says the asker of a build it handed over stopped waiting for it: asked, outcome unknown.
// The build may still run; its outcome, when heard, is the last word.
func (i *Inventory) MarkOutcomeUnknown(ctx context.Context, id, why string) error {
_, err := i.store.Pool().Exec(ctx,
`update build_request set outcome_unknown = $2
where id = $1 and not exists (select 1 from build where build.id = $1)`, id, why)
return err
}
// RequestsNamed is every build request kept whose directory names the module, or whose build said it built
// it, newest first, each with its outcome.
func (i *Inventory) RequestsNamed(ctx context.Context, name string) ([]RequestOutcome, error) {
all, err := i.requests(ctx)
if err != nil {
return nil, err
}
var out []RequestOutcome
for _, a := range all {
if a.Name() == name || a.Module == name {
out = append(out, a)
}
}
return out, nil
}
// RequestedNames is the name of every module a kept build request is expected to register: what `assign`
// offers as a closest name beside the modules registered.
func (i *Inventory) RequestedNames(ctx context.Context) ([]string, error) {
all, err := i.requests(ctx)
if err != nil {
return nil, err
}
seen := map[string]bool{}
var out []string
for _, a := range all {
if n := a.Name(); n != "" && !seen[n] {
seen[n] = true
out = append(out, n)
}
}
return out, nil
}
func (i *Inventory) requests(ctx context.Context) ([]RequestOutcome, error) {
rows, err := i.store.Pool().Query(ctx,
`select a.id, a.repository, a.seat, a.source_path, a.ref, a.commit_hash, a.asked_for,
coalesce(a.not_asked, ''), coalesce(a.outcome_unknown, ''), a.asked_at,
b.id is not null, coalesce(b.failed, ''), coalesce(b.module, ''), coalesce(b.at, a.asked_at)
from build_request a left join build b on b.id = a.id
order by a.asked_at desc, a.id desc`)
if err != nil {
return nil, err
}
defer rows.Close()
var out []RequestOutcome
for rows.Next() {
var a RequestOutcome
if err := rows.Scan(&a.ID, &a.Repository, &a.Seat, &a.Path, &a.Ref, &a.Commit, &a.For, &a.NotAsked,
&a.OutcomeUnknown, &a.At, &a.Heard, &a.Failed, &a.Module, &a.HeardAt); err != nil {
return nil, err
}
out = append(out, a)
}
return out, rows.Err()
}
// The states of a pending assignment.
const (
PendingWaiting = "waiting"
PendingApplying = "applying"
PendingApplied = "applied"
PendingRefused = "refused"
PendingExpired = "expired"
PendingWithdrawn = "withdrawn"
)
// PendingAssignment is an assignment kept until its module is registered (ADR 0261).
type PendingAssignment struct {
ID int64
Node string
Module string
// Build is the build it waits for; Repository and Path where that build is from.
Build, Repository, Path string
Since time.Time
State string
// Note is what ended it, in the words said when it ended.
Note string
Settled *time.Time
// Raised is when an expiry or refusal was raised as a condition; Acknowledged when a person took it
// back; Cleared when its condition was cleared.
Raised, Acknowledged, Cleared *time.Time
}
// Open is whether the pending assignment is still to be made.
func (p PendingAssignment) Open() bool {
return p.State == PendingWaiting || p.State == PendingApplying
}
// ErrAlreadyPending is a second pending assignment of one module to one machine.
var ErrAlreadyPending = errors.New("that assignment is already waiting for its build")
const pendingColumns = `p.id, n.name, p.module, p.build_id, p.repository, p.source_path, p.since, p.state, p.note,
p.settled_at, p.raised_at, p.acknowledged_at, p.cleared_at`
func scanPending(rows pgx.Rows) ([]PendingAssignment, error) {
defer rows.Close()
var out []PendingAssignment
for rows.Next() {
var p PendingAssignment
if err := rows.Scan(&p.ID, &p.Node, &p.Module, &p.Build, &p.Repository, &p.Path, &p.Since, &p.State,
&p.Note, &p.Settled, &p.Raised, &p.Acknowledged, &p.Cleared); err != nil {
return nil, err
}
out = append(out, p)
}
return out, rows.Err()
}
// RecordPending keeps an assignment until its module is registered. One is open per machine and module.
func (i *Inventory) RecordPending(ctx context.Context, p PendingAssignment) (PendingAssignment, error) {
node, err := i.NodeByName(ctx, p.Node)
if err != nil {
return PendingAssignment{}, err
}
err = i.store.Pool().QueryRow(ctx,
`insert into pending_assignment (node, module, build_id, repository, source_path)
values ($1, $2, $3, $4, $5)
on conflict (node, module) where state in ('waiting', 'applying') do nothing
returning id, since`,
node.ID, p.Module, p.Build, p.Repository, strings.Trim(p.Path, "/")).Scan(&p.ID, &p.Since)
if errors.Is(err, pgx.ErrNoRows) {
return PendingAssignment{}, ErrAlreadyPending
}
if err != nil {
return PendingAssignment{}, err
}
p.State = PendingWaiting
return p, nil
}
// RepointPending makes a waiting pending assignment wait for another build. False when it no longer waits.
func (i *Inventory) RepointPending(ctx context.Context, id int64, build, repository, sourcePath string) (bool, error) {
tag, err := i.store.Pool().Exec(ctx,
`update pending_assignment set build_id = $2, repository = $3, source_path = $4
where id = $1 and state = 'waiting'`, id, build, repository, strings.Trim(sourcePath, "/"))
if err != nil {
return false, err
}
return tag.RowsAffected() == 1, nil
}
// Pending is every pending assignment still waiting, oldest first, when endedSince is zero. Otherwise it is
// every open one and every one ended since then.
func (i *Inventory) Pending(ctx context.Context, endedSince time.Time) ([]PendingAssignment, error) {
if endedSince.IsZero() {
rows, err := i.store.Pool().Query(ctx, `select `+pendingColumns+`
from pending_assignment p join node n on n.id = p.node
where p.state = 'waiting' order by p.since, p.id`)
if err != nil {
return nil, err
}
return scanPending(rows)
}
rows, err := i.store.Pool().Query(ctx, `select `+pendingColumns+`
from pending_assignment p join node n on n.id = p.node
where p.state in ('waiting', 'applying') or p.settled_at >= $1
order by p.since, p.id`, endedSince)
if err != nil {
return nil, err
}
return scanPending(rows)
}
// PendingFor is every pending assignment of one module to one machine, newest first.
func (i *Inventory) PendingFor(ctx context.Context, node, module string) ([]PendingAssignment, error) {
rows, err := i.store.Pool().Query(ctx, `select `+pendingColumns+`
from pending_assignment p join node n on n.id = p.node
where n.name = $1 and p.module = $2 order by p.since desc, p.id desc`, node, module)
if err != nil {
return nil, err
}
return scanPending(rows)
}
// Stuck is every pending assignment left applying since before the time given: a controller that claimed
// it and stopped before it said what came of it.
func (i *Inventory) Stuck(ctx context.Context, before time.Time) ([]PendingAssignment, error) {
rows, err := i.store.Pool().Query(ctx, `select `+pendingColumns+`
from pending_assignment p join node n on n.id = p.node
where p.state = 'applying' and p.settled_at < $1 order by p.id`, before)
if err != nil {
return nil, err
}
return scanPending(rows)
}
// ToRaise is every pending assignment that expired or was refused and has not been raised as a condition;
// ToClear every one raised and not cleared.
func (i *Inventory) ToRaise(ctx context.Context) ([]PendingAssignment, error) {
rows, err := i.store.Pool().Query(ctx, `select `+pendingColumns+`
from pending_assignment p join node n on n.id = p.node
where p.state in ('expired', 'refused') and p.raised_at is null order by p.id`)
if err != nil {
return nil, err
}
return scanPending(rows)
}
func (i *Inventory) ToClear(ctx context.Context) ([]PendingAssignment, error) {
rows, err := i.store.Pool().Query(ctx, `select `+pendingColumns+`
from pending_assignment p join node n on n.id = p.node
where p.raised_at is not null and p.cleared_at is null order by p.id`)
if err != nil {
return nil, err
}
return scanPending(rows)
}
// MarkPending stamps a pending assignment raised, acknowledged or cleared, now.
func (i *Inventory) MarkPending(ctx context.Context, id int64, what string) error {
column := map[string]string{"raised": "raised_at", "acknowledged": "acknowledged_at", "cleared": "cleared_at"}[what]
if column == "" {
return errors.New("a pending assignment is marked raised, acknowledged or cleared")
}
_, err := i.store.Pool().Exec(ctx,
`update pending_assignment set `+column+` = coalesce(`+column+`, now()) where id = $1`, id)
return err
}
// ClaimPending takes a waiting pending assignment for making, so nothing else settles or withdraws it while
// it is made. False when it no longer waits: another process took it, or it was withdrawn.
func (i *Inventory) ClaimPending(ctx context.Context, id int64) (bool, error) {
tag, err := i.store.Pool().Exec(ctx,
`update pending_assignment set state = 'applying', settled_at = now() where id = $1 and state = 'waiting'`, id)
if err != nil {
return false, err
}
return tag.RowsAffected() == 1, nil
}
// ReleaseClaim puts a claimed pending assignment back to waiting.
func (i *Inventory) ReleaseClaim(ctx context.Context, id int64) error {
_, err := i.store.Pool().Exec(ctx,
`update pending_assignment set state = 'waiting', settled_at = null where id = $1 and state = 'applying'`, id)
return err
}
// SettlePending ends a pending assignment in the state it is read in (from): waiting for an expiry or a
// withdrawal, applying for the outcome of making it. False when it was no longer in that state.
func (i *Inventory) SettlePending(ctx context.Context, id int64, from, state, note string) (bool, error) {
if state == PendingWaiting || state == PendingApplying {
return false, errors.New("a pending assignment is settled to applied, refused, expired or withdrawn")
}
tag, err := i.store.Pool().Exec(ctx,
`update pending_assignment set state = $3, note = $4, settled_at = now()
where id = $1 and state = $2`, id, from, state, note)
if err != nil {
return false, err
}
return tag.RowsAffected() == 1, nil
}
// ForgetEndedPending deletes ended pending assignments settled more than KeptFor ago, and says how many. One
// raised as a condition is kept until that condition is cleared: deleting it would leave the condition
// with nothing to answer it.
func (i *Inventory) ForgetEndedPending(ctx context.Context, now time.Time) (int64, error) {
tag, err := i.store.Pool().Exec(ctx,
`delete from pending_assignment where state not in ('waiting', 'applying') and settled_at < $1
and (raised_at is null or cleared_at is not null)`,
now.Add(-KeptFor))
if err != nil {
return 0, err
}
return tag.RowsAffected(), nil
}
// PendingKnown says which of these pending assignments are still on record.
func (i *Inventory) PendingKnown(ctx context.Context, ids []int64) (map[int64]bool, error) {
rows, err := i.store.Pool().Query(ctx, `select id from pending_assignment where id = any($1)`, ids)
if err != nil {
return nil, err
}
defer rows.Close()
out := map[int64]bool{}
for rows.Next() {
var id int64
if err := rows.Scan(&id); err != nil {
return nil, err
}
out[id] = true
}
return out, rows.Err()
}
+25
View File
@@ -0,0 +1,25 @@
package inventory
import (
"testing"
"time"
)
// A pending assignment goes with its machine (novox/hq issue 325, review of mesh-controller#150 point 7).
func TestAPendingAssignmentGoesWithItsMachine(t *testing.T) {
inv := ForTest(t)
ctx := t.Context()
if _, err := inv.AddNode(ctx, "leaving"); err != nil {
t.Fatal(err)
}
if _, err := inv.RecordPending(ctx, PendingAssignment{Node: "leaving", Module: "sensors", Build: "b-1"}); err != nil {
t.Fatal(err)
}
if _, err := inv.store.Pool().Exec(ctx, `delete from node where name = 'leaving'`); err != nil {
t.Fatalf("a machine with a pending assignment could not be removed: %v", err)
}
rows, err := inv.Pending(ctx, time.Unix(0, 0))
if err != nil || len(rows) != 0 {
t.Fatalf("a pending assignment outlived its machine: %+v %v", rows, err)
}
}
+9 -4
View File
@@ -76,6 +76,11 @@ func (b *natsBuilds) Ask(ctx context.Context, request BuildRequest) error {
return nil
}
// ErrNotHandedOver is a Submit that failed before the build was handed to the build seat: nothing runs. Any
// other error from Submit came after: the build may still run, and its outcome may still come (novox/hq
// issue 325).
var ErrNotHandedOver = errors.New("the build was not handed over")
func (b *natsBuilds) Submit(ctx context.Context, request BuildRequest,
wait time.Duration) (BuildResult, error) {
@@ -84,23 +89,23 @@ func (b *natsBuilds) Submit(ctx context.Context, request BuildRequest,
// the same event on EVENTS, which the controller records.
outcomes, err := b.js.Conn().SubscribeSync(BuildOutcomeOf(b.role()))
if err != nil {
return BuildResult{}, fmt.Errorf("cannot listen for a build's outcome: %w", err)
return BuildResult{}, fmt.Errorf("%w: cannot listen for a build's outcome: %w", ErrNotHandedOver, err)
}
defer func() { _ = outcomes.Unsubscribe() }()
if err := b.js.Conn().Flush(); err != nil {
return BuildResult{}, err
return BuildResult{}, fmt.Errorf("%w: %w", ErrNotHandedOver, err)
}
body, err := json.Marshal(request)
if err != nil {
return BuildResult{}, err
return BuildResult{}, fmt.Errorf("%w: %w", ErrNotHandedOver, err)
}
// Into the role's work queue and awaited: work the bus never accepted must fail here rather than
// be assumed, because nothing else will ever say so.
publish, cancel := context.WithTimeout(ctx, 30*time.Second)
defer cancel()
if _, err := b.js.Context().Publish(BuildWorkOf(b.role()), body, nats.Context(publish)); err != nil {
return BuildResult{}, fmt.Errorf("cannot submit a build: %w", err)
return BuildResult{}, fmt.Errorf("%w: cannot submit a build: %w", ErrNotHandedOver, err)
}
waiting, cancelWait := context.WithTimeout(ctx, wait)