Compare commits

..
Author SHA1 Message Date
jochen c72f6ca8cb A bundle stands on the toolchain it is compiled in (hq issue 211)
A manifest names its toolchain by language, not in build.on, so the planner did not know a bundle
depends on the module that publishes its toolchain and built the two in one tier: the bundle
against the old toolchain, recorded as built from the new commit. The edge is read from the
manifest, so it holds before any build recorded it, and a toolchain that moves rebuilds every
bundle compiled in it.
2026-10-03 22:19:50 +02:00
43 changed files with 109 additions and 1793 deletions
+2 -23
View File
@@ -148,11 +148,6 @@ func buildFrom(result link.BuildResult) inventory.Build {
// rebuild the graph rather than a list of names. // rebuild the graph rather than a list of names.
Path: result.Path, Path: result.Path,
} }
// When it was asked, which is what orders it against another build of the same module
// (novox/hq 04-ISSUES/219) — not when it was heard.
if asked, ok := link.BuildAskedAt(result.ID); ok {
kept.Asked = asked
}
for _, ref := range result.Against { for _, ref := range result.Against {
kept.Against = append(kept.Against, catalogue.Recorded(ref)) kept.Against = append(kept.Against, catalogue.Recorded(ref))
} }
@@ -414,7 +409,7 @@ func buildOne(ctx context.Context, source buildSource, path, ref string, wait ti
// Correlated by something the control plane makes, not by the module's name: two builds of one // Correlated by something the control plane makes, not by the module's name: two builds of one
// module can be in flight, and the second answer is not the first one's. // module can be in flight, and the second answer is not the first one's.
request := link.BuildRequest{ request := link.BuildRequest{
ID: link.NewBuildID(time.Now()), ID: fmt.Sprintf("%s-%d", "build", time.Now().UnixNano()),
Repository: repository, Repository: repository,
Path: path, Path: path,
Ref: ref, Ref: ref,
@@ -531,31 +526,15 @@ func takeIn(ctx context.Context, inv *inventory.Inventory, result link.BuildResu
BuiltFrom: result.Commit, Head: result.Commit, BuiltFrom: result.Commit, Head: result.Commit,
// What it stood on, so registration can judge a built manifest's base (to-be 38 WP2.4). // What it stood on, so registration can judge a built manifest's base (to-be 38 WP2.4).
Against: kept.Against, Against: kept.Against,
// When it was asked, so an older request heard later does not replace a newer one
// (novox/hq 04-ISSUES/219).
Asked: kept.Asked,
} }
if result.Source != nil && result.Source.Seat != "" { if result.Source != nil && result.Source.Seat != "" {
recorded.Repository, recorded.Seat = result.Source.Repository, result.Source.Seat recorded.Repository, recorded.Seat = result.Source.Repository, result.Source.Seat
} }
// **A build at a commit does not change the branch a module follows** (novox/hq 04-ISSUES/215):
// the commit is built and recorded as what it was built from, and the module keeps following
// what it followed before — the repository's default branch for one new to the catalogue.
if followedBranch(result.Ref) == "" && result.Ref != "" {
recorded.Ref = ""
if was, err := inv.SourceOf(ctx, manifest.Module); err == nil {
recorded.Ref = followedBranch(was.Ref)
}
}
if err := namesNoInstallation(manifest); err != nil { if err := namesNoInstallation(manifest); err != nil {
return manifest, kept, fmt.Errorf("%s built %s (%s), and the mesh does not register it: %w", return manifest, kept, fmt.Errorf("%s built %s (%s), and the mesh does not register it: %w",
result.On, result.Repository, short(result.Commit), err) result.On, result.Repository, short(result.Commit), err)
} }
if err := inv.RegisterModule(ctx, manifest, recorded); err != nil { if err := inv.RegisterModule(ctx, manifest, recorded); err != nil {
if errors.Is(err, inventory.ErrSuperseded) {
return manifest, kept, fmt.Errorf("%s built %s (%s), recorded and not registered: %w",
result.On, manifest.Module, short(result.Commit), err)
}
return manifest, kept, err return manifest, kept, err
} }
return manifest, kept, nil return manifest, kept, nil
@@ -585,7 +564,7 @@ func buildAndShow(ctx context.Context, source buildSource, path, ref string, wai
defer ask.Close() defer ask.Close()
result, err := ask.Submit(ctx, link.BuildRequest{ result, err := ask.Submit(ctx, link.BuildRequest{
ID: link.NewBuildID(time.Now()), ID: fmt.Sprintf("%s-%d", "build", time.Now().UnixNano()),
Repository: repository, Path: path, Ref: ref, Repository: repository, Path: path, Ref: ref,
Held: heldBy(ctx), Seats: seatBases(ctx), Held: heldBy(ctx), Seats: seatBases(ctx),
}, wait) }, wait)
-87
View File
@@ -2,12 +2,9 @@ package main
import ( import (
"encoding/json" "encoding/json"
"errors"
"strings" "strings"
"testing" "testing"
"time"
"github.com/novox/mesh-controller/internal/inventory"
"github.com/novox/mesh-controller/internal/link" "github.com/novox/mesh-controller/internal/link"
) )
@@ -66,87 +63,3 @@ func TestABuildHeardIsRecordedAndRegistered(t *testing.T) {
t.Fatalf("a failure is said in the builder's words: %v", err) t.Fatalf("a failure is said in the builder's words: %v", err)
} }
} }
// novox/hq 04-ISSUES/215: a build asked at a commit is recorded as built from that commit, and the
// module keeps following the branch it followed — a new one, the default branch.
func TestABuildAtACommitKeepsTheBranchTheModuleFollows(t *testing.T) {
open := aMesh(t)
ctx := t.Context()
manifest, _ := json.Marshal(map[string]any{"module": "unifi", "version": "1"})
result := func(id, ref, commit string) link.BuildResult {
return link.BuildResult{ID: id, Repository: "http://forge.internal:20000/novox/mesh-catalog.git",
Path: "modules/unifi", Ref: ref, On: "anchor", Commit: commit, Manifest: manifest,
Source: &link.SourceOnSeat{Seat: "git", Repository: "novox/mesh-catalog"}}
}
if _, _, err := takeIn(ctx, open.inventory, result("b-1", "main", "1111111aaaa")); err != nil {
t.Fatal(err)
}
if _, _, err := takeIn(ctx, open.inventory, result("b-2", "9c97a8a", "9c97a8a1d2c3")); err != nil {
t.Fatal(err)
}
src, err := open.inventory.SourceOf(ctx, "unifi")
if err != nil {
t.Fatal(err)
}
if src.Ref != "main" || src.BuiltFrom != "9c97a8a1d2c3" {
t.Errorf("after a build at a commit the module follows %q, built from %q; want main, 9c97a8a1d2c3", src.Ref, src.BuiltFrom)
}
// One new to the catalogue, first built at a commit, follows the default branch.
other, _ := json.Marshal(map[string]any{"module": "letta", "version": "1"})
r := result("b-3", "deadbeef", "deadbeefcafe")
r.Manifest, r.Path = other, "modules/letta"
if _, _, err := takeIn(ctx, open.inventory, r); err != nil {
t.Fatal(err)
}
if src, _ := open.inventory.SourceOf(ctx, "letta"); src.Ref != "" {
t.Errorf("a module first built at a commit follows %q, want the default branch", src.Ref)
}
}
// novox/hq 04-ISSUES/219: an older request heard after a newer one is recorded and not registered,
// so a push sends what the newer request built.
func TestAnOlderBuildHeardLaterDoesNotReplaceTheNewer(t *testing.T) {
open := aMesh(t)
ctx := t.Context()
older := time.Date(2026, 10, 3, 21, 33, 45, 0, time.UTC)
newer := time.Date(2026, 10, 3, 21, 51, 57, 0, time.UTC)
result := func(asked time.Time, image string) link.BuildResult {
manifest, _ := json.Marshal(map[string]any{"module": "postgres", "version": image})
return link.BuildResult{ID: link.NewBuildID(asked), Repository: "http://forge.internal:20000/novox/mesh-catalog.git",
Path: "modules/postgres", Ref: "main", On: "anchor", Commit: "efff5415", Manifest: manifest,
Source: &link.SourceOnSeat{Seat: "git", Repository: "novox/mesh-catalog"}}
}
if _, _, err := takeIn(ctx, open.inventory, result(newer, "4bcd5f73")); err != nil {
t.Fatal(err)
}
_, _, err := takeIn(ctx, open.inventory, result(older, "0ab07fa9"))
if !errors.Is(err, inventory.ErrSuperseded) {
t.Fatalf("the older request's outcome was taken in as current: %v", err)
}
shelf, err := open.inventory.Catalogue(ctx)
if err != nil {
t.Fatal(err)
}
if got := shelf["postgres"].Version; got != "4bcd5f73" {
t.Errorf("postgres is %q; want the newer request's 4bcd5f73", got)
}
if builds, _ := open.inventory.Builds(ctx, "postgres", 5); len(builds) != 2 {
t.Errorf("the late build was not recorded: %v", builds)
}
}
func TestABuildIDSaysWhenItWasAsked(t *testing.T) {
at := time.Date(2026, 10, 3, 21, 51, 57, 392539762, time.UTC)
if got, ok := link.BuildAskedAt(link.NewBuildID(at)); !ok || !got.Equal(at) {
t.Errorf("read back %v %v; want %v", got, ok, at)
}
if got, ok := link.BuildAskedAt("build-1791064317392539762"); !ok || got.Format(time.TimeOnly) != "21:51:57" {
t.Errorf("the incident's id reads as %v %v", got, ok)
}
for _, id := range []string{"b-1", "build-2", "build-", "build-x", ""} {
if _, ok := link.BuildAskedAt(id); ok {
t.Errorf("%q read as a request time", id)
}
}
}
+3 -13
View File
@@ -8,7 +8,6 @@ package main
import ( import (
"context" "context"
"errors"
"flag" "flag"
"fmt" "fmt"
"os" "os"
@@ -254,34 +253,25 @@ func parseAround(set *flag.FlagSet, args []string) ([]string, error) {
// became of a build nobody was watching. // became of a build nobody was watching.
func (b builds) Built(ctx context.Context, result link.BuildResult) error { func (b builds) Built(ctx context.Context, result link.BuildResult) error {
manifest, _, err := takeIn(ctx, b.inv, result) manifest, _, err := takeIn(ctx, b.inv, result)
// 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)
switch { switch {
case err != nil && result.Failed != "": case err != nil && result.Failed != "":
fmt.Printf("%s: %v\n", result.ID, err) fmt.Printf("%s: %v\n", result.ID, err)
if result.Module != "" { if result.Module != "" {
planBuilt(ctx, b.open, result.Module, result.Commit, result.Failed, asked) planBuilt(ctx, b.open, result.Module, result.Commit, result.Failed)
} else { } else {
planFailedBuild(ctx, b.open, result) planFailedBuild(ctx, b.open, result)
} }
return nil return nil
case errors.Is(err, inventory.ErrSuperseded):
// Not a failure: the module is already at what a later request built. A plan that asked
// before that later request is answered by it; one that asked after it ignores this.
fmt.Printf("%s: %v\n", result.ID, err)
planBuilt(ctx, b.open, manifest.Module, result.Commit, "", asked)
return nil
case err != nil: case err != nil:
fmt.Printf("%s: heard and recorded, and not registered: %v\n", result.ID, err) fmt.Printf("%s: heard and recorded, and not registered: %v\n", result.ID, err)
if manifest.Module != "" { if manifest.Module != "" {
planBuilt(ctx, b.open, manifest.Module, result.Commit, err.Error(), asked) planBuilt(ctx, b.open, manifest.Module, result.Commit, err.Error())
} }
return nil return nil
} }
fmt.Printf("%s: %s %s registered, built on %s from %s\n", fmt.Printf("%s: %s %s registered, built on %s from %s\n",
result.ID, manifest.Module, manifest.Version, result.On, short(result.Commit)) result.ID, manifest.Module, manifest.Version, result.On, short(result.Commit))
saysWhenThePolicyActs(ctx, b.inv, manifest.Module) saysWhenThePolicyActs(ctx, b.inv, manifest.Module)
planBuilt(ctx, b.open, manifest.Module, result.Commit, "", asked) planBuilt(ctx, b.open, manifest.Module, result.Commit, "")
return nil return nil
} }
-6
View File
@@ -266,12 +266,6 @@ func TestTheResolverIsToldEveryMachineOnTheNetworkAndToldAgainWhenOneLeaves(t *t
if _, err := assign(ctx, open, "anchor", "dnsmasq"); err != nil { if _, err := assign(ctx, open, "anchor", "dnsmasq"); err != nil {
t.Fatal(err) t.Fatal(err)
} }
// Its bus credential, as assigning issues it where the bus is reachable (novox/hq issue 203):
// no bus is known to this test, so it is minted here, or composing refuses the placeholder.
if _, err := open.inventory.MintBusPassword(ctx, inventory.BusUser{
Username: "anchor.dnsmasq", Kind: inventory.BusModule, Node: "anchor", Module: "dnsmasq"}); err != nil {
t.Fatal(err)
}
zones := func() string { zones := func() string {
t.Helper() t.Helper()
for _, r := range composed(t, open, "anchor").Resources { for _, r := range composed(t, open, "anchor").Resources {
-24
View File
@@ -211,27 +211,3 @@ func TestWhatAHandedOverModuleRecordsAboutItsSource(t *testing.T) {
} }
} }
} }
// novox/hq 04-ISSUES/215: a module once built at a commit still follows its branch — a merge into it
// matches the module, and a plan re-asks the branch, not the old commit.
func TestAModuleBuiltAtACommitStillFollowsItsBranch(t *testing.T) {
m := link.SourceMoved{Owner: "novox", Repo: "mesh-catalog", Base: "main"}
pinned := inventory.Source{Repository: "novox/mesh-catalog", Seat: "git", Ref: "9c97a8a"}
if !sourceIs(pinned, m) {
t.Error("a module whose record names a commit is left out of a merge into its branch")
}
full := inventory.Source{Repository: "novox/mesh-catalog", Seat: "git", Ref: "9c97a8a1d2c3b4a5f60718293a4b5c6d7e8f9012"}
if !sourceIs(full, m) {
t.Error("a full commit hash is read as a branch")
}
if got := followedBranch("9c97a8a"); got != "" {
t.Errorf("a plan would re-ask the old commit %q", got)
}
if got := followedBranch("release"); got != "release" {
t.Errorf("a branch is not followed as named: %q", got)
}
// A module that follows another branch is still not this merge's.
if sourceIs(inventory.Source{Repository: "novox/mesh-catalog", Seat: "git", Ref: "release"}, m) {
t.Error("a module following another branch was matched")
}
}
@@ -1,47 +0,0 @@
package main
import (
"testing"
"time"
"github.com/novox/mesh-controller/internal/inventory"
)
// novox/hq issue 213: for the moment a machine hands its controller over, the container and the
// process both run the plan timer on one store. Only the one holding the plans moves them; the other
// leaves them alone, and moves them once they are let go.
func TestAControllerLeavesThePlansToTheOneHoldingThem(t *testing.T) {
open := aMesh(t)
ctx := t.Context()
now := time.Now().UTC()
// Every tier done: the next step is the plan's last, and needs nothing but the store.
plan := inventory.Plan{ID: "plan-213", Repository: "r", Commit: "abc", Created: now, Updated: now,
State: inventory.PlanRolling, Tier: 1, Tiers: [][]string{{"app"}},
Modules: map[string]*inventory.PlanModule{"app": {State: "built"}}}
if err := open.inventory.SavePlan(ctx, plan); err != nil {
t.Fatal(err)
}
// The other controller: its own connections to the same store, holding the plans.
other, err := inventory.Open(ctx)
if err != nil {
t.Fatal(err)
}
t.Cleanup(other.Close)
release, err := other.HoldPlans(ctx, false)
if err != nil {
t.Fatal(err)
}
t.Cleanup(release) // before the close above: a pool waits for a connection still held
advancePlans(ctx, open)
if p, err := open.inventory.PlanByID(ctx, "plan-213"); err != nil || !p.Open() {
t.Fatalf("a controller moved a plan another held: %+v %v", p, err)
}
release()
advancePlans(ctx, open)
if p, err := open.inventory.PlanByID(ctx, "plan-213"); err != nil || p.State != inventory.PlanDone {
t.Fatalf("the plan did not move once it was let go: %+v %v", p, err)
}
}
+4 -122
View File
@@ -2,7 +2,6 @@ package main
import ( import (
"context" "context"
"errors"
"flag" "flag"
"fmt" "fmt"
"sort" "sort"
@@ -273,8 +272,7 @@ func askTier(ctx context.Context, inv *inventory.Inventory, p *inventory.Plan) e
} }
source := buildSource{Repository: e.Source.Repository, Seat: e.Source.Seat} source := buildSource{Repository: e.Source.Repository, Seat: e.Source.Seat}
fmt.Printf(" tier %d: ", p.Tier) fmt.Printf(" tier %d: ", p.Tier)
// The branch it follows, never a commit a build once named (novox/hq 04-ISSUES/215). if err := buildOne(ctx, source, e.Source.Path, e.Source.Ref, 0); err != nil {
if err := buildOne(ctx, source, e.Source.Path, followedBranch(e.Source.Ref), 0); err != nil {
state.State = "failed" state.State = "failed"
state.Why = err.Error() state.Why = err.Error()
p.State = inventory.PlanFailed p.State = inventory.PlanFailed
@@ -289,23 +287,8 @@ func askTier(ctx context.Context, inv *inventory.Inventory, p *inventory.Plan) e
// planBuilt marks a module built (or failed) in every open plan whose current tier holds it, and // planBuilt marks a module built (or failed) in every open plan whose current tier holds it, and
// advances what that completes. Called from the daemon's take-in of every outcome. // advances what that completes. Called from the daemon's take-in of every outcome.
// func planBuilt(ctx context.Context, open *stores, module, commit, failed string) {
// **Only a build asked at or after the plan's ask is its outcome** (novox/hq 04-ISSUES/219). Two
// plans a few minutes apart both ask for a module; the earlier plan's build, finishing late, is not
// the later plan's answer — it stood on the bases from before the later plan's merge, and taking it
// would send machines, and the next tier, what the later merge replaced. asked is zero when the
// build's request time is not known, and such an outcome is taken as before.
func planBuilt(ctx context.Context, open *stores, module, commit, failed string, asked time.Time) {
inv := open.inventory inv := open.inventory
// One controller works the plans at a time (novox/hq issue 213); an outcome waits its turn rather
// than write over what the holder is about to save. Not taken, it is still in the build records,
// which the holder settles the plan from (issue 214).
release, err := inv.HoldPlans(ctx, true)
if err != nil {
fmt.Printf("plans: %s's outcome is left to the build records: %v\n", module, err)
return
}
defer release()
plans, err := inv.OpenPlans(ctx) plans, err := inv.OpenPlans(ctx)
if err != nil { if err != nil {
fmt.Printf("plans: cannot read them: %v\n", err) fmt.Printf("plans: cannot read them: %v\n", err)
@@ -331,9 +314,6 @@ func planBuilt(ctx context.Context, open *stores, module, commit, failed string,
state = &inventory.PlanModule{} state = &inventory.PlanModule{}
p.Modules[module] = state p.Modules[module] = state
} }
if askedBefore(asked, state.AskedAt) {
continue
}
if failed != "" { if failed != "" {
state.State = "failed" state.State = "failed"
state.Why = failed state.Why = failed
@@ -352,30 +332,13 @@ func planBuilt(ctx context.Context, open *stores, module, commit, failed string,
fmt.Printf("%s: %s; the tiers after it are not asked\n", p.ID, p.Note) fmt.Printf("%s: %s; the tiers after it are not asked\n", p.ID, p.Note)
} }
} }
advanceHeld(ctx, open) advancePlans(ctx, open)
} }
// advancePlans moves every open plan as far as the facts allow: a tier whose modules are all built // advancePlans moves every open plan as far as the facts allow: a tier whose modules are all built
// and whose gates are applied gives way to the next; the last tier done is the plan done. Called // and whose gates are applied gives way to the next; the last tier done is the plan done. Called
// after every outcome and on a timer, so a plan waiting on a machine's report moves when it comes. // after every outcome and on a timer, so a plan waiting on a machine's report moves when it comes.
//
// **One controller at a time** (novox/hq issue 213). A plan is read, changed and saved whole; two
// controllers — the old and the new while a machine hands its controller over — would each ask a
// tier the other had just asked. Taken without waiting: whoever holds the plans is moving them.
func advancePlans(ctx context.Context, open *stores) { func advancePlans(ctx context.Context, open *stores) {
release, err := open.inventory.HoldPlans(ctx, false)
if err != nil {
if !errors.Is(err, inventory.ErrPlansBusy) {
fmt.Printf("plans: cannot hold them: %v\n", err)
}
return
}
defer release()
advanceHeld(ctx, open)
}
// advanceHeld is advancePlans for a caller already holding the plans.
func advanceHeld(ctx context.Context, open *stores) {
inv := open.inventory inv := open.inventory
plans, err := inv.OpenPlans(ctx) plans, err := inv.OpenPlans(ctx)
if err != nil { if err != nil {
@@ -436,24 +399,6 @@ func advanceOnce(ctx context.Context, open *stores, p *inventory.Plan,
} }
return true, nil return true, nil
} }
// **Asked: settle from the build records first** (novox/hq 04-ISSUES/214). An outcome is taken
// in by whichever controller hears it, and a merge to the controller's own repository replaces
// the controller in its first tier: the build that produced the new one is recorded, and the
// plan never hears it. The record is the fact; a build recorded after the ask is that tier's
// outcome, whoever was listening.
recorded := map[string][]inventory.Build{}
for _, m := range tier {
if s := p.Modules[m]; s != nil && s.State == "asked" {
builds, err := inv.Builds(ctx, m, 5)
if err != nil {
return false, err
}
recorded[m] = builds
}
}
if settleFromRecords(p, tier, recorded) {
return true, nil
}
// Asked: wait for every build. // Asked: wait for every build.
var latest time.Time var latest time.Time
for _, m := range tier { for _, m := range tier {
@@ -585,8 +530,7 @@ func planFailedBuild(ctx context.Context, open *stores, result link.BuildResult)
} }
for _, e := range entries { for _, e := range entries {
if repositoryMatches(e.Source.Repository, result.Repository) && e.Source.Path == result.Path { if repositoryMatches(e.Source.Repository, result.Repository) && e.Source.Path == result.Path {
asked, _ := link.BuildAskedAt(result.ID) planBuilt(ctx, open, e.Manifest.Module, result.Commit, result.Failed)
planBuilt(ctx, open, e.Manifest.Module, result.Commit, result.Failed, asked)
return return
} }
} }
@@ -703,19 +647,6 @@ func plansCommand(ctx context.Context, args []string) error {
} }
p.State = inventory.PlanFailed p.State = inventory.PlanFailed
p.Note = "stopped by hand at tier " + fmt.Sprint(p.Tier) p.Note = "stopped by hand at tier " + fmt.Sprint(p.Tier)
release, err := inv.HoldPlans(ctx, true)
if err != nil {
return err
}
defer release()
if p, err = inv.PlanByID(ctx, positionals[1]); err != nil {
return err
}
if !p.Open() {
return fmt.Errorf("%s is already %s", p.ID, p.State)
}
p.State = inventory.PlanFailed
p.Note = "stopped by hand at tier " + fmt.Sprint(p.Tier)
if err := inv.SavePlan(ctx, p); err != nil { if err := inv.SavePlan(ctx, p); err != nil {
return err return err
} }
@@ -824,52 +755,3 @@ func splitList(s string) []string {
} }
return out return out
} }
// settleFromRecords marks every module of the tier still `asked` built — or failed — from a build
// recorded after it was asked, and says whether it changed anything (novox/hq 04-ISSUES/214).
// Newest first, as Builds answers: the first record after the ask is the outcome of that ask.
func settleFromRecords(p *inventory.Plan, tier []string, recorded map[string][]inventory.Build) bool {
changed := false
for _, m := range tier {
s := p.Modules[m]
if s == nil || s.State != "asked" || s.AskedAt == nil {
continue
}
var outcome *inventory.Build
for i := range recorded[m] {
b := recorded[m][i]
if b.At.Before(*s.AskedAt) {
break
}
// Recorded after the ask and asked before it: an earlier ask's late outcome, not this
// one's (novox/hq 04-ISSUES/219).
if askedBefore(b.Asked, s.AskedAt) {
continue
}
outcome = &b
}
if outcome == nil {
continue
}
at := outcome.At
if outcome.Worked() {
s.State = "built"
s.BuiltAt = &at
s.Commit = outcome.Commit
} else {
s.State = "failed"
s.Why = outcome.Failed
p.State = inventory.PlanFailed
p.Note = fmt.Sprintf("%s failed to build in tier %d", m, p.Tier)
}
fmt.Printf("%s: %s settled from the build records as %s (%s)\n", p.ID, m, s.State, outcome.ID)
changed = true
}
return changed
}
// askedBefore is whether a build asked at asked was asked before a plan asked for its module — and
// so is not that plan's outcome (novox/hq 04-ISSUES/219). False when either time is not known.
func askedBefore(asked time.Time, planAsked *time.Time) bool {
return !asked.IsZero() && planAsked != nil && asked.Before(*planAsked)
}
-55
View File
@@ -126,58 +126,3 @@ func TestABundleIsPlannedAfterTheToolchainItIsCompiledIn(t *testing.T) {
t.Fatalf("the toolchain, then the bundle: %v", p.Tiers) t.Fatalf("the toolchain, then the bundle: %v", p.Tiers)
} }
} }
// novox/hq 04-ISSUES/214: a plan whose build outcome was recorded while no controller followed it —
// the controller rebuilding itself — settles from the build records instead of waiting for ever.
func TestAPlanSettlesAnAskedBuildFromTheRecords(t *testing.T) {
asked := time.Date(2026, 10, 3, 19, 20, 0, 0, time.UTC)
p := inventory.Plan{ID: "plan-1", Tiers: [][]string{{"mesh-controller", "builder"}, {"route-proxy"}},
Modules: map[string]*inventory.PlanModule{
"mesh-controller": {State: "asked", AskedAt: &asked},
"builder": {State: "asked", AskedAt: &asked},
}}
records := map[string][]inventory.Build{
// Newest first, as Builds answers: the build after the ask is the outcome.
"mesh-controller": {
{ID: "build-2", Commit: "2ebbb799", At: asked.Add(4 * time.Minute)},
{ID: "build-1", Commit: "06ea2168", At: asked.Add(-10 * time.Minute)},
},
// Only a build from before the ask: not this ask's outcome.
"builder": {{ID: "build-0", Commit: "06ea2168", At: asked.Add(-time.Hour)}},
}
if !settleFromRecords(&p, p.Tiers[0], records) {
t.Fatal("nothing settled, though the controller's build is recorded after the ask")
}
if s := p.Modules["mesh-controller"]; s.State != "built" || s.Commit != "2ebbb799" || s.BuiltAt == nil {
t.Errorf("the controller's ask is %+v, want built from 2ebbb799", s)
}
if s := p.Modules["builder"]; s.State != "asked" {
t.Errorf("an ask with no record after it was settled: %+v", s)
}
// novox/hq 04-ISSUES/219: a build recorded after the ask but asked before it — an earlier
// plan's late outcome — is not this ask's, built or failed.
r := inventory.Plan{ID: "plan-3", Tiers: [][]string{{"postgres"}},
Modules: map[string]*inventory.PlanModule{"postgres": {State: "asked", AskedAt: &asked}}}
late := map[string][]inventory.Build{"postgres": {
{ID: "build-old", Commit: "efff5415", Asked: asked.Add(-18 * time.Minute), At: asked.Add(12 * time.Minute)},
}}
if settleFromRecords(&r, r.Tiers[0], late) || r.Modules["postgres"].State != "asked" {
t.Errorf("an earlier ask's late outcome settled this ask: %+v", r.Modules["postgres"])
}
// Newest heard first: the earlier ask's late outcome, then this ask's own, heard before it.
late["postgres"] = append(late["postgres"], inventory.Build{ID: "build-mine", Commit: "4bcd5f73",
Asked: asked.Add(time.Second), At: asked.Add(5 * time.Minute)})
if !settleFromRecords(&r, r.Tiers[0], late) || r.Modules["postgres"].State != "built" ||
r.Modules["postgres"].Commit != "4bcd5f73" {
t.Errorf("this ask's own outcome, heard before the earlier ask's, did not settle it: %+v", r.Modules["postgres"])
}
// A failure recorded after the ask fails the plan, as hearing it would have.
q := inventory.Plan{ID: "plan-2", Tiers: [][]string{{"x"}},
Modules: map[string]*inventory.PlanModule{"x": {State: "asked", AskedAt: &asked}}}
settleFromRecords(&q, q.Tiers[0], map[string][]inventory.Build{"x": {{ID: "b", Failed: "no", At: asked.Add(time.Minute)}}})
if q.State != inventory.PlanFailed || q.Modules["x"].State != "failed" {
t.Errorf("a recorded failure did not fail the plan: %+v %+v", q, q.Modules["x"])
}
}
+3 -6
View File
@@ -402,22 +402,19 @@ func seatAnnouncement(handlers map[string]link.ToolHandler) micro.Info {
var endpoints []micro.EndpointInfo var endpoints []micro.EndpointInfo
for _, verb := range verbs { for _, verb := range verbs {
schema, _ := json.Marshal(about[verb].Input) schema, _ := json.Marshal(about[verb].Input)
// The same shape every tool runtime announces in (node-tools' announce package): the name is
// `<seat>__<verb>`, as the protocol's characters allow; the metadata is what identifies it.
endpoints = append(endpoints, micro.EndpointInfo{ endpoints = append(endpoints, micro.EndpointInfo{
Name: catalogue.ControllerSeatName + "__" + verb, Name: verb,
Subject: link.SeatToolSubject(catalogue.ControllerSeatName, verb), Subject: link.SeatToolSubject(catalogue.ControllerSeatName, verb),
QueueGroup: "seat." + catalogue.ControllerSeatName, QueueGroup: "seat." + catalogue.ControllerSeatName,
Metadata: map[string]string{ Metadata: map[string]string{
"kind": "seat", "module": catalogue.ControllerSeatName, "tool": verb,
"seat": catalogue.ControllerSeatName, "scope": "mesh", "interchangeable": "false",
"description": about[verb].Description, "schema": string(schema), "description": about[verb].Description, "schema": string(schema),
"seat": catalogue.ControllerSeatName, "scope": "mesh",
}, },
}) })
} }
return micro.Info{ return micro.Info{
ServiceIdentity: micro.ServiceIdentity{ ServiceIdentity: micro.ServiceIdentity{
Name: catalogue.ControllerSeatName, ID: "controller", Version: "0.1.0", Name: catalogue.ControllerSeatName, ID: "controller", Version: "1.0.0",
Metadata: map[string]string{"seat": catalogue.ControllerSeatName, "scope": "mesh"}, Metadata: map[string]string{"seat": catalogue.ControllerSeatName, "scope": "mesh"},
}, },
Description: "the mesh's own verbs, answered by the holder of the mesh-controller seat", Description: "the mesh's own verbs, answered by the holder of the mesh-controller seat",
+3 -7
View File
@@ -281,14 +281,10 @@ func TestTheControllerAnnouncesTheVerbsItServes(t *testing.T) {
t.Fatalf("%d endpoints announced for %d verbs served", len(info.Endpoints), len(handlers)) t.Fatalf("%d endpoints announced for %d verbs served", len(info.Endpoints), len(handlers))
} }
for _, e := range info.Endpoints { for _, e := range info.Endpoints {
verb := e.Metadata["tool"] if _, served := handlers[e.Name]; !served {
if _, served := handlers[verb]; !served || e.Name != catalogue.ControllerSeatName+"__"+verb { t.Errorf("%s is announced and not served", e.Name)
t.Errorf("%s (%s) is announced and not served under that name", e.Name, verb)
} }
if e.Metadata["kind"] != "seat" || e.Metadata["seat"] != catalogue.ControllerSeatName { if e.Subject != link.SeatToolSubject(catalogue.ControllerSeatName, e.Name) || e.QueueGroup != "seat."+catalogue.ControllerSeatName {
t.Errorf("%s is not announced as the seat's verb: %v", e.Name, e.Metadata)
}
if e.Subject != link.SeatToolSubject(catalogue.ControllerSeatName, verb) || e.QueueGroup != "seat."+catalogue.ControllerSeatName {
t.Errorf("%s is announced on %s/%s, not where it is served", e.Name, e.Subject, e.QueueGroup) t.Errorf("%s is announced on %s/%s, not where it is served", e.Name, e.Subject, e.QueueGroup)
} }
if e.Metadata["description"] == "" || e.Metadata["schema"] == "" || e.Metadata["scope"] != "mesh" { if e.Metadata["description"] == "" || e.Metadata["schema"] == "" || e.Metadata["scope"] != "mesh" {
+1 -34
View File
@@ -5,7 +5,6 @@ import (
"errors" "errors"
"flag" "flag"
"fmt" "fmt"
"regexp"
"strings" "strings"
"time" "time"
@@ -284,14 +283,6 @@ func (f following) SourceMoved(ctx context.Context, m link.SourceMoved) error {
if isHistory(m.MergedAt, lastLookAt(entries, m)) { if isHistory(m.MergedAt, lastLookAt(entries, m)) {
packaging = nil packaging = nil
} }
// Said, never silent (novox/hq 04-ISSUES/215): a module built from this repository that follows
// another branch is not part of this merge, and whoever is waiting for its change should read why.
for _, e := range entries {
if sameRepository(e.Source.Repository, m) && !sourceIs(e.Source, m) {
fmt.Printf(" %s is built from %s/%s and follows %s, not %s; this merge leaves it out\n",
e.Manifest.Module, m.Owner, m.Repo, e.Source.Ref, m.Base)
}
}
touched := whatTheMergeTouched(from, entries, m) touched := whatTheMergeTouched(from, entries, m)
for _, e := range touched { for _, e := range touched {
if err := inv.SourceMoved(ctx, e.Manifest.Module, m.Commit); err != nil { if err := inv.SourceMoved(ctx, e.Manifest.Module, m.Commit); err != nil {
@@ -315,13 +306,6 @@ func (f following) SourceMoved(ctx context.Context, m link.SourceMoved) error {
for _, e := range moved { for _, e := range moved {
movedNames = append(movedNames, e.Manifest.Module) movedNames = append(movedNames, e.Manifest.Module)
} }
// Written and its first tier asked as one act on the plans (novox/hq issue 213): a timer on
// another controller reading it between the two would ask the tier again.
release, err := inv.HoldPlans(ctx, true)
if err != nil {
return notNow(err)
}
defer release()
plan := planOfMerge(m, movedNames, edges) plan := planOfMerge(m, movedNames, edges)
if hasCycle(plan.Tiers, edges) { if hasCycle(plan.Tiers, edges) {
fmt.Printf(" the last tier depends on itself: %s — built together, in no order\n", fmt.Printf(" the last tier depends on itself: %s — built together, in no order\n",
@@ -361,24 +345,7 @@ func sourceIs(s inventory.Source, m link.SourceMoved) bool {
if !sameRepository(s.Repository, m) { if !sameRepository(s.Repository, m) {
return false return false
} }
ref := followedBranch(s.Ref) return s.Ref == "" || s.Ref == m.Base
return ref == "" || ref == m.Base
}
// commitRef is a ref that names a commit rather than a branch: what `build --ref <commit>` asks for.
var commitRef = regexp.MustCompile(`^[0-9a-f]{7,40}$`)
// followedBranch is the branch a recorded ref means a module follows (novox/hq 04-ISSUES/215). **A
// commit is never a branch to follow.** A build asked at a commit — to try one, or to pin it during a
// fix — recorded that commit as the module's ref; every merge after it then failed to match the
// module, its plan left it out without saying so, and every plan that rebuilt it asked for that same
// old commit again. A commit recorded so is read as the repository's default branch, which is what
// the module followed before it; a branch is followed as named.
func followedBranch(ref string) string {
if commitRef.MatchString(strings.TrimSpace(ref)) {
return ""
}
return ref
} }
// sameRepository is whether a recorded repository is the one a merge names, in either spelling it // sameRepository is whether a recorded repository is the one a merge names, in either spelling it
-29
View File
@@ -1,29 +0,0 @@
package main
import (
"crypto/tls"
"net/http"
"net/http/httptest"
"net/url"
"testing"
)
// A backend behind the proxy learns the client used TLS and which name it asked for, so the addresses
// it writes into its own pages are the ones a client can use (2026-10-03: a forge's Go import tag
// named an http clone URL, and Go refused the module path).
func TestABackendIsToldTheRequestWasHTTPSAndForWhichName(t *testing.T) {
var proto, host, fwdHost, fwdFor string
backend := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
proto, host, fwdHost, fwdFor = r.Header.Get("X-Forwarded-Proto"), r.Host, r.Header.Get("X-Forwarded-Host"), r.Header.Get("X-Forwarded-For")
}))
defer backend.Close()
where, _ := url.Parse(backend.URL)
req := httptest.NewRequest(http.MethodGet, "https://git.example.org/novox/mesh-sdk/go?go-get=1", nil)
req.TLS = &tls.ConnectionState{}
req.Host = "git.example.org"
req.RemoteAddr = "192.0.2.7:51000"
towards(where).ServeHTTP(httptest.NewRecorder(), req)
if proto != "https" || fwdHost != "git.example.org" || host != "git.example.org" || fwdFor != "192.0.2.7" {
t.Errorf("the backend was told proto=%q host=%q forwarded-host=%q for=%q", proto, host, fwdHost, fwdFor)
}
}
+1 -15
View File
@@ -273,7 +273,7 @@ func (t *table) set(routes map[string][]rule, public map[string]bool) {
log.Printf("route %s points at %q, which is not a URL: %v", host, r.target, err) log.Printf("route %s points at %q, which is not a URL: %v", host, r.target, err)
continue continue
} }
r.to = towards(where) r.to = httputil.NewSingleHostReverseProxy(where)
if r.insecure { if r.insecure {
r.to.Transport = &http.Transport{TLSClientConfig: &tls.Config{InsecureSkipVerify: true}} r.to.Transport = &http.Transport{TLSClientConfig: &tls.Config{InsecureSkipVerify: true}}
} }
@@ -1060,17 +1060,3 @@ func asPort(v any) (int, bool) {
} }
return 0, false return 0, false
} }
// towards proxies to one backend and tells it what the client asked: **X-Forwarded-Proto, -Host and
// -For**, set from the request this proxy received. A backend that builds its own addresses — a forge
// writing its clone URL into a page, a login redirect — otherwise sees the plain HTTP hop from this
// proxy and writes `http://`, though every client reached it over TLS: Go refused the forge's module
// path for exactly that on 2026-10-03, its import tag naming an http clone URL.
// The standard library's NewSingleHostReverseProxy sets only X-Forwarded-For.
func towards(where *url.URL) *httputil.ReverseProxy {
return &httputil.ReverseProxy{Rewrite: func(pr *httputil.ProxyRequest) {
pr.SetURL(where)
pr.Out.Host = pr.In.Host
pr.SetXForwarded()
}}
}
+5 -25
View File
@@ -465,28 +465,10 @@ func PermissionsFor(p Principal) (Permissions, error) {
pub = append(pub, invoked...) pub = append(pub, invoked...)
// It says what it serves and may ask what answers (novox/hq ADR 0197): the runtime answers // It says what it serves and may ask what answers (novox/hq ADR 0197): the runtime answers
// discovery for each module and seat it carries, and the console it is asks the bus. // discovery for each module and seat it carries, and the console it is asks the bus.
// One service per runtime process, named for the runtime: the bus lets a principal answer each sub = append(sub, announcing(serves...)...)
// request once, so the runtime announces everything it carries under its own name.
sub = append(sub, announcing(append([]string{RuntimeModule}, serves...)...)...)
pub = append(pub, discovering()...) pub = append(pub, discovering()...)
// **And it consumes for the modules it carries** (novox/hq ADR 0198, which changes ADR 0175's // Nothing about consumers: it consumes nothing. A module's reactions to events are its
// "it consumes nothing"): a module's long-running code is a bundle this runtime launches, and // own long-lived process, which ADR 0175 leaves where it is; what moves here is tools.
// the runtime is its bus — it reads the module's own durable consumer and acknowledges what
// the module's code took. Exactly the grants the module's own principal has for that consumer,
// on its name and no other's: asking about it, pulling from it, acknowledging it. The
// consumer is still the controller's to make, from the module's own principal.
for _, d := range p.Carries {
own := Principal{Kind: KindModule, Node: p.Node, Module: d.Module, Emits: d.Emits,
Consumes: d.Consumes, Serves: d.Serves, Holds: d.Holds, Uses: d.Uses, Watches: d.Watches}
if _, consumes := ConsumerFor(own); !consumes {
continue
}
stream, durable := consumerStream(own), consumerDurable(own)
pub = append(pub,
"$JS.API.CONSUMER.INFO."+stream+"."+durable,
"$JS.API.CONSUMER.MSG.NEXT."+stream+"."+durable,
"$JS.ACK."+stream+"."+durable+".>")
}
sub = unique(sub) sub = unique(sub)
pub = unique(pub) pub = unique(pub)
} }
@@ -811,14 +793,12 @@ func invokedSubjects(invokes []string) ([]string, error) {
// protocol's discovery (novox/hq ADR 0197): the questions asked of every service, and those asked of // protocol's discovery (novox/hq ADR 0197): the questions asked of every service, and those asked of
// each name it serves — its own and no other's, so it cannot answer for a service it is not. // each name it serves — its own and no other's, so it cannot answer for a service it is not.
func announcing(names ...string) []string { func announcing(names ...string) []string {
out := []string{"$SRV.PING", "$SRV.INFO", "$SRV.STATS"} out := []string{"$SRV.PING", "$SRV.INFO"}
for _, n := range names { for _, n := range names {
if !safeSubject.MatchString(n) { if !safeSubject.MatchString(n) {
continue continue
} }
for _, verb := range []string{"PING", "INFO", "STATS"} { out = append(out, "$SRV.PING."+n, "$SRV.PING."+n+".>", "$SRV.INFO."+n, "$SRV.INFO."+n+".>")
out = append(out, "$SRV."+verb+"."+n, "$SRV."+verb+"."+n+".>")
}
} }
return out return out
} }
+8 -19
View File
@@ -378,7 +378,7 @@ func TestAModulePullsItsOwnConsumerAndNoOthers(t *testing.T) {
// their tools (novox/hq ADR 0175): every carried module's tool namespace, every held seat's verbs // their tools (novox/hq ADR 0175): every carried module's tool namespace, every held seat's verbs
// on this node, every module's membership on this node, and a call to anything. Nothing it // on this node, every module's membership on this node, and a call to anything. Nothing it
// consumes, because it reacts to nothing. // consumes, because it reacts to nothing.
func TestTheRuntimeServesTheUnionAndConsumesForItsModules(t *testing.T) { func TestTheRuntimeServesTheUnionAndConsumesNothing(t *testing.T) {
filter := Seat{Name: "node-packet-filter", Scope: "node", Serves: []string{"rules", "reload"}} filter := Seat{Name: "node-packet-filter", Scope: "node", Serves: []string{"rules", "reload"}}
p := Principal{Kind: KindNodeTools, Node: "anchor", Module: RuntimeModule, Carries: []Declared{ p := Principal{Kind: KindNodeTools, Node: "anchor", Module: RuntimeModule, Carries: []Declared{
{Module: "nftables", Holds: []Seat{filter}, Serves: []string{"firewall_rules"}}, {Module: "nftables", Holds: []Seat{filter}, Serves: []string{"firewall_rules"}},
@@ -408,33 +408,22 @@ func TestTheRuntimeServesTheUnionAndConsumesForItsModules(t *testing.T) {
t.Errorf("the runtime may not publish %s: %v", want, perms.Publish) t.Errorf("the runtime may not publish %s: %v", want, perms.Publish)
} }
} }
// It reads the consumer of every carried module that consumes — that module's, by its name, as // Nothing of what a carried module consumes, and no consumer of its own to ack.
// the module's own principal could (novox/hq ADR 0198) — and of no module that consumes nothing. for _, s := range perms.Subscribe {
for _, want := range []string{ if strings.Contains(s, ".event.") || strings.HasPrefix(s, "_DELIVER.") {
"$JS.API.CONSUMER.INFO.EVENTS.anchor_zsh", t.Errorf("the runtime was granted a delivery it has no consumer for: %s", s)
"$JS.API.CONSUMER.MSG.NEXT.EVENTS.anchor_zsh",
"$JS.ACK.EVENTS.anchor_zsh.>",
} {
if !contains(perms.Publish, want) {
t.Errorf("the runtime may not read zsh's consumer: %s missing from %v", want, perms.Publish)
} }
} }
for _, s := range perms.Publish { for _, s := range perms.Publish {
if (strings.HasPrefix(s, "$JS.ACK.") || strings.Contains(s, "CONSUMER")) && !strings.Contains(s, "anchor_zsh") { if strings.HasPrefix(s, "$JS.ACK.") || strings.Contains(s, "CONSUMER") {
t.Errorf("the runtime was granted a consumer no carried module of it consumes on: %s", s) t.Errorf("the runtime was granted a consumer's subject and has no consumer: %s", s)
}
}
// It pulls; nothing is pushed to it, and it subscribes no event subject directly.
for _, s := range perms.Subscribe {
if strings.Contains(s, ".event.") || strings.HasPrefix(s, "_DELIVER.") {
t.Errorf("the runtime was granted a delivery: %s", s)
} }
} }
if !perms.AllowResponses { if !perms.AllowResponses {
t.Error("the runtime answers what it is asked, and may not reply") t.Error("the runtime answers what it is asked, and may not reply")
} }
if _, needed := ConsumerFor(p); needed { if _, needed := ConsumerFor(p); needed {
t.Error("a consumer would be made for the runtime itself; it reads its modules' consumers, never one of its own") t.Error("a consumer would be made for the runtime, which consumes nothing")
} }
// Each subject once in each list: the file is read as the mesh's authority model. One subject may // Each subject once in each list: the file is read as the mesh's authority model. One subject may
// stand in both — the runtime answers discovery on `$SRV.INFO` and, as the console, asks it // stand in both — the runtime answers discovery on `$SRV.INFO` and, as the console, asks it
+4 -4
View File
@@ -25,7 +25,7 @@ accounts {
users = [ users = [
{ user: "controller", password: "$2a$11$cccccccccccccccccccccc", permissions: { { user: "controller", password: "$2a$11$cccccccccccccccccccccc", permissions: {
publish: { allow: ["$JS.ACK.CONTROL.controller.>", "$JS.ACK.EVENTS.controller.>", "$JS.API.>", "_INBOX.enrol.>", "mesh.assignment.>", "mesh.control.>", "mesh.mod.*.tool.>", "mesh.node.>", "mesh.seat.mesh-build-machine.accept.>", "mesh.seat.mesh-controller.event.applied", "mesh.seat.mesh-controller.event.built-before", "mesh.seat.mesh-controller.event.refused", "mesh.seat.node-build-agent.accept.>"] } publish: { allow: ["$JS.ACK.CONTROL.controller.>", "$JS.ACK.EVENTS.controller.>", "$JS.API.>", "_INBOX.enrol.>", "mesh.assignment.>", "mesh.control.>", "mesh.mod.*.tool.>", "mesh.node.>", "mesh.seat.mesh-build-machine.accept.>", "mesh.seat.mesh-controller.event.applied", "mesh.seat.mesh-controller.event.built-before", "mesh.seat.mesh-controller.event.refused", "mesh.seat.node-build-agent.accept.>"] }
subscribe: { allow: ["$JS.API.>", "$SRV.INFO", "$SRV.INFO.mesh-controller", "$SRV.INFO.mesh-controller.>", "$SRV.PING", "$SRV.PING.mesh-controller", "$SRV.PING.mesh-controller.>", "$SRV.STATS", "$SRV.STATS.mesh-controller", "$SRV.STATS.mesh-controller.>", "_DELIVER.controller", "_DELIVER.controller.>", "_INBOX.controller.>", "mesh.control.>", "mesh.mod.gitea.event.pull.merged", "mesh.mod.mesh-catalog.event.catching-up", "mesh.mod.mesh-catalog.event.upgraded", "mesh.seat.mesh-build-machine.event.built", "mesh.seat.mesh-controller.tool.>", "mesh.seat.node-build-agent.event.built"] } subscribe: { allow: ["$JS.API.>", "$SRV.INFO", "$SRV.INFO.mesh-controller", "$SRV.INFO.mesh-controller.>", "$SRV.PING", "$SRV.PING.mesh-controller", "$SRV.PING.mesh-controller.>", "_DELIVER.controller", "_DELIVER.controller.>", "_INBOX.controller.>", "mesh.control.>", "mesh.mod.gitea.event.pull.merged", "mesh.mod.mesh-catalog.event.catching-up", "mesh.mod.mesh-catalog.event.upgraded", "mesh.seat.mesh-build-machine.event.built", "mesh.seat.mesh-controller.tool.>", "mesh.seat.node-build-agent.event.built"] }
allow_responses: { max: 1, ttl: "1m" } allow_responses: { max: 1, ttl: "1m" }
} } } }
{ user: "enrol.one", password: "$2a$11$eeeeeeeeeeeeeeeeeeeeee", permissions: { { user: "enrol.one", password: "$2a$11$eeeeeeeeeeeeeeeeeeeeee", permissions: {
@@ -38,17 +38,17 @@ accounts {
} } } }
{ user: "one.telegram", password: "$2a$11$tttttttttttttttttttttt", permissions: { { user: "one.telegram", password: "$2a$11$tttttttttttttttttttttt", permissions: {
publish: { allow: ["$JS.ACK.EVENTS.one_telegram.>", "$JS.ACK.SEAT_TELEGRAM_SENDER.SEAT_TELEGRAM_SENDER_worker.>", "$JS.API.CONSUMER.INFO.EVENTS.one_telegram", "$JS.API.CONSUMER.INFO.SEAT_TELEGRAM_SENDER.SEAT_TELEGRAM_SENDER_worker", "$JS.API.CONSUMER.MSG.NEXT.EVENTS.one_telegram", "$JS.API.CONSUMER.MSG.NEXT.SEAT_TELEGRAM_SENDER.SEAT_TELEGRAM_SENDER_worker", "$JS.API.DIRECT.GET.ASSIGNMENTS.mesh.assignment.one.telegram", "mesh.seat.telegram-sender.event.delivered", "mesh.seat.telegram-sender.event.failed"] } publish: { allow: ["$JS.ACK.EVENTS.one_telegram.>", "$JS.ACK.SEAT_TELEGRAM_SENDER.SEAT_TELEGRAM_SENDER_worker.>", "$JS.API.CONSUMER.INFO.EVENTS.one_telegram", "$JS.API.CONSUMER.INFO.SEAT_TELEGRAM_SENDER.SEAT_TELEGRAM_SENDER_worker", "$JS.API.CONSUMER.MSG.NEXT.EVENTS.one_telegram", "$JS.API.CONSUMER.MSG.NEXT.SEAT_TELEGRAM_SENDER.SEAT_TELEGRAM_SENDER_worker", "$JS.API.DIRECT.GET.ASSIGNMENTS.mesh.assignment.one.telegram", "mesh.seat.telegram-sender.event.delivered", "mesh.seat.telegram-sender.event.failed"] }
subscribe: { allow: ["$SRV.INFO", "$SRV.INFO.telegram", "$SRV.INFO.telegram.>", "$SRV.PING", "$SRV.PING.telegram", "$SRV.PING.telegram.>", "$SRV.STATS", "$SRV.STATS.telegram", "$SRV.STATS.telegram.>", "_INBOX.one.telegram.>", "mesh.assignment.one.telegram", "mesh.mod.telegram.tool.>", "mesh.seat.telegram-sender.accept.send"] } subscribe: { allow: ["$SRV.INFO", "$SRV.INFO.telegram", "$SRV.INFO.telegram.>", "$SRV.PING", "$SRV.PING.telegram", "$SRV.PING.telegram.>", "_INBOX.one.telegram.>", "mesh.assignment.one.telegram", "mesh.mod.telegram.tool.>", "mesh.seat.telegram-sender.accept.send"] }
allow_responses: { max: 1, ttl: "1m" } allow_responses: { max: 1, ttl: "1m" }
} } } }
{ user: "two.audit", password: "$2a$11$aaaaaaaaaaaaaaaaaaaaaa", permissions: { { user: "two.audit", password: "$2a$11$aaaaaaaaaaaaaaaaaaaaaa", permissions: {
publish: { allow: ["$JS.ACK.EVENTS.two_audit.>", "$JS.API.CONSUMER.INFO.EVENTS.two_audit", "$JS.API.CONSUMER.MSG.NEXT.EVENTS.two_audit", "$JS.API.DIRECT.GET.ASSIGNMENTS.mesh.assignment.two.audit"] } publish: { allow: ["$JS.ACK.EVENTS.two_audit.>", "$JS.API.CONSUMER.INFO.EVENTS.two_audit", "$JS.API.CONSUMER.MSG.NEXT.EVENTS.two_audit", "$JS.API.DIRECT.GET.ASSIGNMENTS.mesh.assignment.two.audit"] }
subscribe: { allow: ["$SRV.INFO", "$SRV.INFO.audit", "$SRV.INFO.audit.>", "$SRV.PING", "$SRV.PING.audit", "$SRV.PING.audit.>", "$SRV.STATS", "$SRV.STATS.audit", "$SRV.STATS.audit.>", "_INBOX.two.audit.>", "mesh.assignment.two.audit", "mesh.mod.audit.tool.>", "mesh.mod.shop.event.order.placed"] } subscribe: { allow: ["$SRV.INFO", "$SRV.INFO.audit", "$SRV.INFO.audit.>", "$SRV.PING", "$SRV.PING.audit", "$SRV.PING.audit.>", "_INBOX.two.audit.>", "mesh.assignment.two.audit", "mesh.mod.audit.tool.>", "mesh.mod.shop.event.order.placed"] }
allow_responses: { max: 1, ttl: "1m" } allow_responses: { max: 1, ttl: "1m" }
} } } }
{ user: "two.shop", password: "$2a$11$ssssssssssssssssssssss", permissions: { { user: "two.shop", password: "$2a$11$ssssssssssssssssssssss", permissions: {
publish: { allow: ["$JS.ACK.EVENTS.two_shop.>", "$JS.API.CONSUMER.INFO.EVENTS.two_shop", "$JS.API.CONSUMER.MSG.NEXT.EVENTS.two_shop", "$JS.API.DIRECT.GET.ASSIGNMENTS.mesh.assignment.two.shop", "mesh.mod.shop.event.order.placed", "mesh.seat.telegram-sender.accept.send"] } publish: { allow: ["$JS.ACK.EVENTS.two_shop.>", "$JS.API.CONSUMER.INFO.EVENTS.two_shop", "$JS.API.CONSUMER.MSG.NEXT.EVENTS.two_shop", "$JS.API.DIRECT.GET.ASSIGNMENTS.mesh.assignment.two.shop", "mesh.mod.shop.event.order.placed", "mesh.seat.telegram-sender.accept.send"] }
subscribe: { allow: ["$SRV.INFO", "$SRV.INFO.shop", "$SRV.INFO.shop.>", "$SRV.PING", "$SRV.PING.shop", "$SRV.PING.shop.>", "$SRV.STATS", "$SRV.STATS.shop", "$SRV.STATS.shop.>", "_INBOX.two.shop.>", "mesh.assignment.two.shop", "mesh.mod.shop.tool.>"] } subscribe: { allow: ["$SRV.INFO", "$SRV.INFO.shop", "$SRV.INFO.shop.>", "$SRV.PING", "$SRV.PING.shop", "$SRV.PING.shop.>", "_INBOX.two.shop.>", "mesh.assignment.two.shop", "mesh.mod.shop.tool.>"] }
allow_responses: { max: 1, ttl: "1m" } allow_responses: { max: 1, ttl: "1m" }
} } } }
] ]
+1 -99
View File
@@ -593,12 +593,6 @@ func one(ctx context.Context, run Runner, publish Publisher,
if err != nil { if err != nil {
return catalogue.Built{}, fmt.Errorf("%s: writing %s's launchers failed: %w", module, a.Name, err) return catalogue.Built{}, fmt.Errorf("%s: writing %s's launchers failed: %w", module, a.Name, err)
} }
if chain.Bundler != "" {
say("bundle", "bundling each entrypoint into one file")
if compiled, err = bundled(ctx, run, tree, chain, base, a, launchers); err != nil {
return catalogue.Built{}, fmt.Errorf("%s: bundling %s failed: %w", module, a.Name, err)
}
}
say("bundle", "compiled, packing") say("bundle", "compiled, packing")
body, err := pack(compiled) body, err := pack(compiled)
if err != nil { if err != nil {
@@ -980,7 +974,7 @@ func compile(ctx context.Context, run Runner, tree string, chain Toolchain,
if _, err := run(ctx, tree, "docker", invocation...); err != nil { if _, err := run(ctx, tree, "docker", invocation...); err != nil {
return "", err return "", err
} }
if chain.Dependencies != "" && chain.Bundler == "" { if chain.Dependencies != "" {
// **What the bundle runs with, from the image it was compiled in** (Toolchain.Dependencies). // **What the bundle runs with, from the image it was compiled in** (Toolchain.Dependencies).
// A second run in the same image rather than a shell wrapped around the compiler: the // A second run in the same image rather than a shell wrapped around the compiler: the
// compile line stays a plain command a reader can run by hand, and the copy is one more // compile line stays a plain command a reader can run by hand, and the copy is one more
@@ -1230,95 +1224,3 @@ func writeLaunchers(root string, chain Toolchain, a catalogue.Artifact) (map[str
} }
return out, nil return out, nil
} }
// bundledSuffix is where a bundle's one-file output is written, beside what the compiler wrote.
const bundledSuffix = ".bundled"
// bundled makes every entrypoint and every launcher of a compiled bundle ONE file, in the toolchain
// image's bundler, and answers the directory to pack (novox/hq ADR 0193).
//
// **What a launched bundle runs is what it imports, and nothing else.** Every served bundle is its
// own process, so it carries its own copy of the SDK and its own dependencies inlined — the
// toolchain's whole node_modules no longer travels in every bundle. An entrypoint a process runs by
// name (`node daemon/index.js`) is bundled in place under its own name; a launcher keeps its name
// and its first line, and stays executable. A package the bundler cannot inline is named by the
// artifact (`external`), kept as an import, and only then is the toolchain's runtime directory
// copied beside the files. CommonJS inlined into an ES module still finds `require`.
func bundled(ctx context.Context, run Runner, tree string, chain Toolchain, base string,
a catalogue.Artifact, launchers map[string]string) (string, error) {
const within = "/app/modules/module"
out, final := Out(a.Name), Out(a.Name)+bundledSuffix
if err := os.RemoveAll(filepath.Join(tree, final)); err != nil {
return "", err
}
if err := os.MkdirAll(filepath.Join(tree, final), 0o755); err != nil {
return "", err
}
common := []string{"--bundle", "--platform=node", "--format=esm", "--target=node22",
"--outbase=" + out, "--outdir=" + final, "--log-level=warning",
"--banner:js=import { createRequire as __meshRequire } from 'node:module'; const require = __meshRequire(import.meta.url);"}
for _, x := range a.External {
common = append(common, "--external:"+x)
}
var plain []string
for _, e := range a.Entrypoints {
if strings.HasSuffix(e, ".js") {
plain = append(plain, out+"/"+e)
}
}
var launch []string
for _, l := range sortedValues(launchers) {
launch = append(launch, out+"/"+l)
}
// Refused by name in an image that predates the bundler, as the dependencies copy is: a bundle
// packed without it would carry nothing it imports. Run as itself: npm installs esbuild's native
// binary in place of its script, which `node` cannot run.
guard := `test -x "$0" || { echo "the toolchain image carries no bundler at $0: it predates one-file bundles, rebuild mesh-tools first" >&2; exit 1; }; exec "$0" "$@"`
step := func(entries []string, extra ...string) error {
if len(entries) == 0 {
return nil
}
invocation := []string{"run", "--rm", "--volume", tree + ":" + within, "--workdir", within, base,
"sh", "-c", guard, chain.Bundler}
invocation = append(invocation, entries...)
invocation = append(invocation, common...)
invocation = append(invocation, extra...)
_, err := run(ctx, tree, "docker", invocation...)
return err
}
if err := step(plain); err != nil {
return "", err
}
if err := step(launch, "--out-extension:.js=.mjs"); err != nil {
return "", err
}
// Plain `.js` output is an ES module; said once, as the runtime directory used to say it.
if err := os.WriteFile(filepath.Join(tree, final, "package.json"), []byte(`{"type":"module","private":true}`+"\n"), 0o644); err != nil {
return "", err
}
for _, l := range launchers {
path := filepath.Join(tree, final, filepath.FromSlash(l))
if _, err := os.Stat(path); err == nil {
if err := os.Chmod(path, 0o755); err != nil {
return "", err
}
}
}
if len(a.External) > 0 && chain.Dependencies != "" {
copying := []string{"run", "--rm", "--volume", tree + ":" + within, "--workdir", within, base,
"sh", "-c", `cp -a "$0/node_modules" "$1/"`, chain.Dependencies, final}
if _, err := run(ctx, tree, "docker", copying...); err != nil {
return "", fmt.Errorf("copying the packages %s keeps external: %w", a.Name, err)
}
}
return filepath.Join(tree, final), nil
}
func sortedValues(m map[string]string) []string {
out := make([]string, 0, len(m))
for _, v := range m {
out = append(out, v)
}
sort.Strings(out)
return out
}
+14 -38
View File
@@ -83,49 +83,25 @@ func TestABundleIsCompiledAndPackedWithNoDockerfile(t *testing.T) {
t.Fatalf("the bundle was not pinned: %v", got.Manifest.Resources[0]) t.Fatalf("the bundle was not pinned: %v", got.Manifest.Resources[0])
} }
// **One file per entrypoint and launcher, in the toolchain's bundler** (novox/hq ADR 0193). A // **And what it runs with, from the image it was compiled in** (novox/hq to-be 38 WP3). A
// second run in the same toolchain image bundles each into the artifact's bundled output, the SDK // second run in the same toolchain image copies the toolchain's runtime directory — the
// inlined, refusing by name in an image that predates the bundler; and the toolchain's // `"type": "module"` package.json and the pruned node_modules — into the output's root, and
// node_modules is no longer copied into a bundle that keeps nothing external. // refuses by name when the image carries none rather than packing a bundle that starts nowhere.
var bundling []string var copied string
for _, line := range r.ran { for _, line := range r.ran {
if strings.HasPrefix(line, "docker run") && strings.Contains(line, "esbuild") { if strings.HasPrefix(line, "docker run") && strings.Contains(line, "/app/runtime") {
bundling = append(bundling, line) copied = line
} }
} }
if len(bundling) != 2 { if copied == "" {
t.Fatalf("want one bundling run for the entrypoints and one for the launchers:\n%s", strings.Join(r.ran, "\n")) t.Fatalf("the bundle's dependencies were not copied in after the compile:\n%s", strings.Join(r.ran, "\n"))
} }
for _, want := range []string{"mesh-tools/build@sha256:", "predates one-file bundles", "--bundle", "--format=esm", if !strings.Contains(copied, "mesh-tools/build@sha256:") || !strings.Contains(copied, "predates") ||
"--platform=node", "--outdir=" + Out("code") + ".bundled", Out("code") + "/index.js"} { !strings.Contains(copied, Out("code")) {
if !strings.Contains(bundling[0], want) { t.Fatalf("the copy does not run in the same toolchain, refuse an older image by name, or land in the artifact's output: %s", copied)
t.Errorf("the entrypoints' bundling lacks %q: %s", want, bundling[0])
}
} }
if !strings.Contains(bundling[1], Out("code")+"/index.serve.mjs") || !strings.Contains(bundling[1], "--out-extension:.js=.mjs") { if strings.Index(strings.Join(r.ran, "\n"), "--outDir") > strings.Index(strings.Join(r.ran, "\n"), "/app/runtime") {
t.Errorf("the launcher is not bundled under its own name: %s", bundling[1]) t.Fatal("the dependencies were copied before the compile wrote its output")
}
if strings.Contains(strings.Join(r.ran, "\n"), "/app/runtime") {
t.Errorf("the toolchain's node_modules was copied into a bundle that keeps nothing external:\n%s", strings.Join(r.ran, "\n"))
}
if strings.Index(strings.Join(r.ran, "\n"), "--outDir") > strings.Index(strings.Join(r.ran, "\n"), "esbuild") {
t.Fatal("the bundler ran before the compile wrote its output")
}
}
// A bundle naming packages it keeps external is bundled with them as imports, and carries the
// toolchain's node_modules for them — the one case it still does.
func TestABundleKeepingAPackageExternalCarriesTheToolchainsModules(t *testing.T) {
manifest := strings.Replace(aBundle, `"entrypoints":["index.js"]`, `"entrypoints":["index.js"],"external":["sharp"]`, 1)
r, workspace := aRepository(t, manifest, map[string]string{"index.ts": "console.log(1)"})
held := map[string]string{"mesh-tools/build": "registry.invalid/mesh-tools/build@sha256:" + strings.Repeat("b", 64)}
if _, err := Build(context.Background(), compiling{r}.run, r,
"https://forge.invalid/greeter.git", "", "", workspace, held, Npmrc{}, GitCredential{}, nil); err != nil {
t.Fatal(err)
}
all := strings.Join(r.ran, "\n")
if !strings.Contains(all, "--external:sharp") || !strings.Contains(all, "/app/runtime") {
t.Errorf("an external package was not kept as an import with the toolchain's modules beside it:\n%s", all)
} }
} }
-20
View File
@@ -200,23 +200,3 @@ func TestWhatABuildReadIsTheRepositoriesItsRecipesName(t *testing.T) {
t.Fatal("a module whose recipes name no other repository read one") t.Fatal("a module whose recipes name no other repository read one")
} }
} }
// novox/hq 04-ISSUES/212: a toolchain stands on the SDK's published package, and is built with the
// exact version the mesh published — an argument that changes when the SDK does, so a rebuild after
// a release never reuses an install of the version before it.
func TestAPackageTheMeshPublishedIsPassedByItsExactVersion(t *testing.T) {
manifest := catalogue.Manifest{
Module: "mesh-tools",
Build: &catalogue.Build{
On: []catalogue.BuildsOn{{Arg: "MESH_SDK", Module: "mesh-sdk", Artifact: "lib"}},
},
}
held := map[string]string{"mesh-sdk/lib": "@novox/mesh-sdk@0.1.6"}
args, resolved, err := standingOn(context.Background(), manifest, held, noMirror)
if err != nil {
t.Fatal(err)
}
if fmt.Sprint(args) != "[--build-arg MESH_SDK=@novox/mesh-sdk@0.1.6]" || fmt.Sprint(resolved) != "[@novox/mesh-sdk@0.1.6]" {
t.Errorf("the package was passed as %v, recorded as %v", args, resolved)
}
}
-10
View File
@@ -73,16 +73,7 @@ type Toolchain struct {
// //
// A toolchain image without the directory fails the build by name rather than packing a bundle // A toolchain image without the directory fails the build by name rather than packing a bundle
// that starts nowhere: the image predates this and must be rebuilt first. // that starts nowhere: the image predates this and must be rebuilt first.
//
// *Since the bundler (below):* copied only for a bundle that names packages it keeps external,
// which cannot be inlined; a bundle with none carries no node_modules at all.
Dependencies string Dependencies string
// Bundler is the bundler inside the toolchain image that makes each compiled entrypoint and each
// launcher ONE self-contained file (novox/hq ADR 0193): every served bundle is its own process
// now, so each carries its own copy of what it imports — the SDK included — and nothing else.
// A bundle shrinks from the toolchain's whole node_modules to the code it runs. Empty for a
// language whose build is already one file.
Bundler string
// SystemStamp is the variable this language's linker fills with the artifact's declared system, // SystemStamp is the variable this language's linker fills with the artifact's declared system,
// for a language whose binaries are pinned to one at link time (novox/hq ADR 0005). // for a language whose binaries are pinned to one at link time (novox/hq ADR 0005).
// //
@@ -148,7 +139,6 @@ var toolchains = []Toolchain{
Unit: UnitSources, Unit: UnitSources,
SourceExt: ".ts", SourceExt: ".ts",
Dependencies: "/app/runtime", Dependencies: "/app/runtime",
Bundler: "/app/node_modules/esbuild/bin/esbuild",
}, },
{ {
Language: "go", Language: "go",
@@ -1,166 +0,0 @@
package catalogue
import (
"encoding/json"
"strings"
"testing"
)
// novox/hq issue 213 (ADR 0188 §1, §3): a module's own Go service is a bundle the host runs as a
// process, not an image. Each test holds one thing that had to change in the composer for the
// controller to be declared that way.
const aServiceDigest = "sha256:dddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddd"
// aServiceModule is the controller's shape in miniature: it answers tools of its own, its code is a
// Go bundle a process runs as an account it declares, its secrets belong to that account, it
// prepares its state, and its process replaces the container it used to run as.
func aServiceModule(t *testing.T) Manifest {
t.Helper()
raw := `{
"module": "svc", "version": "1", "tools": ["status"], "prepares": true,
"own-secrets": {"store": "${dir:state}/store"},
"secrets-owner": "svc",
"resources": [
{"id": "state", "type": "directory", "mode": "0700", "place": "mesh", "owner": "svc"},
{"id": "service", "type": "process", "name": "svc", "artifact": "code",
"run": ["./svc", "serve"], "user": "svc", "replaces": ["server"],
"env": {"SVC_STORE_FILE": "${dir:state}/store", "SVC_STORE_PORT": "${seat:mesh-store:5432}"}},
{"id": "account", "type": "user", "name": "svc", "shell": "/usr/bin/nologin", "home": "/var/lib/svc"}
],
"build": {"artifacts": [{"name": "code", "kind": "bundle", "language": "go", "system": "arch",
"from": "cmd/svc", "binary": "svc"}]}
}`
m, err := ParseManifest([]byte(raw))
if err != nil {
t.Fatalf("the service's manifest is refused: %v", err)
}
resolved, err := m.Resolve([]Built{{Name: "code", Kind: ArtifactBundle,
Reference: ArtifactStoreScheme + "svc/code@" + aServiceDigest, Digest: aServiceDigest}})
if err != nil {
t.Fatal(err)
}
return resolved
}
func composeTheService(t *testing.T, with Rendering) []map[string]any {
t.Helper()
with.Needed = map[string]map[string]string{"svc": {"store": "sealed-store"}}
with.ArtifactStore = "anchor.internal:5100"
out, err := Resolution{Node: "anchor", Modules: []Manifest{aServiceModule(t)}}.Declaration(with)
if err != nil {
t.Fatalf("the service does not compose: %v", err)
}
return out
}
func indexOf(out []map[string]any, id string) int {
for i, r := range out {
if r["id"] == id {
return i
}
}
return -1
}
// A module that declares tools has every bundle served by the node's runtime unless it says
// otherwise — and the controller declares the verbs it answers as tools. Its service bundle is run
// by its own process; launched a second time by the runtime it would be a second controller
// pretending to be an MCP server.
func TestABundleItsOwnProcessRunsIsNotServedByTheRuntime(t *testing.T) {
m := aServiceModule(t)
if len(m.Bundles) != 1 {
t.Fatalf("the service's bundle was not kept: %+v", m.Bundles)
}
if loads := m.Bundles[0].Loads; len(loads) != 0 {
t.Fatalf("the runtime would launch the service's own bundle as tools: %v", loads)
}
// And a bundle no resource runs still is served, as a module declaring tools always had it.
tools := Manifest{Module: "t", Version: "1", Tools: []string{"x"},
Build: &Build{Artifacts: []Artifact{{Name: "tools", Kind: ArtifactBundle, Language: "go",
System: "arch", From: "cmd/t"}}}}
resolved, err := tools.Resolve([]Built{{Name: "tools", Kind: ArtifactBundle,
Reference: ArtifactStoreScheme + "t/tools@" + aServiceDigest, Digest: aServiceDigest}})
if err != nil {
t.Fatal(err)
}
if loads := resolved.Bundles[0].Loads; len(loads) != 1 || loads[0] != "t" {
t.Fatalf("a tools bundle nothing runs is no longer served: %v", loads)
}
}
// The account is created before anything is given to it. Its secrets are mesh-computed and so
// placed before the module's own resources; given to a user the machine did not have yet, they were
// refused on the first apply and the process started without them.
func TestAModulesAccountComesBeforeWhatBelongsToIt(t *testing.T) {
out := composeTheService(t, Rendering{})
account, secret := indexOf(out, "svc.account"), indexOf(out, "svc."+NeedID("store"))
if account < 0 || secret < 0 {
t.Fatalf("the account or the secret is missing: %v", out)
}
if account > secret {
t.Fatalf("the secret owned by svc is written before svc exists: account at %d, secret at %d",
account, secret)
}
if owner := out[secret]["owner"]; owner != "svc" {
t.Errorf("the secret belongs to %v, not the account its process runs as", owner)
}
}
// The process is the module's program; its preparation is the same program asked to prepare, as a
// step before it — with the same account and environment, and handing nothing over.
func TestAProcessIsPreparedByItsOwnProgram(t *testing.T) {
out := composeTheService(t, Rendering{})
step, process := indexOf(out, "svc.service-prepare"), indexOf(out, "svc.service")
if step < 0 || process < 0 || step > process {
t.Fatalf("the preparation is not a step before the process (%d, %d): %v", step, process, out)
}
s := out[step]
if s["type"] != "process" || s["run-once"] != true || s["name"] != "svc-prepare" {
t.Errorf("the preparation is not a run-once process: %v", s)
}
if run, _ := json.Marshal(s["run"]); string(run) != `["./svc","prepare"]` {
t.Errorf("the preparation runs %s", run)
}
if s["user"] != "svc" || s["source"] != out[process]["source"] {
t.Errorf("the preparation does not run the same bundle as the same account: %v", s)
}
if env, _ := s["env"].(map[string]any); env["SVC_STORE_FILE"] == nil {
t.Errorf("the preparation is not given the process's environment: %v", s["env"])
}
if _, has := s["replaces"]; has {
t.Errorf("the preparation would hand over what the process replaces: %v", s)
}
if _, has := s["args"]; has {
t.Errorf("the preparation carries a container's args: %v", s)
}
}
// What the process replaces is named as the host recorded it, `<module>.<id>`; unprefixed, the host
// matches nothing and removes the container first, as before.
func TestWhatAProcessReplacesIsNamedAsTheHostRecordedIt(t *testing.T) {
out := composeTheService(t, Rendering{})
p := out[indexOf(out, "svc.service")]
if got, _ := json.Marshal(p["replaces"]); string(got) != `["svc.server"]` {
t.Fatalf("the process replaces %s", got)
}
}
func TestWhatReplacesMayNameIsRefusedNearItsAuthor(t *testing.T) {
for what, resource := range map[string]string{
"a container": `{"id":"c","type":"container","name":"c","image":"x@` + aServiceDigest + `","replaces":["old"]}`,
"a step": `{"id":"p","type":"process","name":"p","run":["./p"],"run-once":true,"replaces":["old"]}`,
"something declared": `{"id":"p","type":"process","name":"p","run":["./p"],"replaces":["p"]}`,
"another module's": `{"id":"p","type":"process","name":"p","run":["./p"],"replaces":["other.old"]}`,
"not a list": `{"id":"p","type":"process","name":"p","run":["./p"],"replaces":"old"}`,
} {
raw := `{"module":"m","version":"1","resources":[` + resource + `]}`
if _, err := ParseManifest([]byte(raw)); err == nil || !strings.Contains(err.Error(), "replace") {
t.Errorf("replaces on %s was accepted: %v", what, err)
}
}
ok := `{"module":"m","version":"1","resources":[{"id":"p","type":"process","name":"p","run":["./p"],"replaces":["old"]}]}`
if _, err := ParseManifest([]byte(ok)); err != nil {
t.Errorf("a process replacing what its module no longer declares was refused: %v", err)
}
}
+1 -60
View File
@@ -77,7 +77,6 @@ func (m Manifest) Resolve(built []Built) (Manifest, error) {
// build compare equal. // build compare equal.
out.Bundles = nil out.Bundles = nil
if m.Build != nil { if m.Build != nil {
run := runByAResource(m)
for _, a := range m.Build.Artifacts { for _, a := range m.Build.Artifacts {
if a.Kind != ArtifactBundle { if a.Kind != ArtifactBundle {
continue continue
@@ -85,20 +84,9 @@ func (m Manifest) Resolve(built []Built) (Manifest, error) {
made := by[a.Name] made := by[a.Name]
// What the runtime loads: what the artifact said, else every entrypoint of a module // What the runtime loads: what the artifact said, else every entrypoint of a module
// that declares tools, else nothing (the field's own rule; see Artifact.Loads). // that declares tools, else nothing (the field's own rule; see Artifact.Loads).
//
// **Never, unasked, a bundle one of the module's own resources runs** (novox/hq issue 213).
// A process the host runs is the module's service, not its tools: the controller declares
// the verbs it answers as `tools` and serves them itself, and its bundle would otherwise
// have been launched a second time by the node's runtime, as an MCP child it is not.
loads := append([]string(nil), a.Loads...) loads := append([]string(nil), a.Loads...)
if a.Loads == nil && len(m.Tools) > 0 && !run[a.Name] { if a.Loads == nil && len(m.Tools) > 0 {
loads = append([]string(nil), a.Entrypoints...) loads = append([]string(nil), a.Entrypoints...)
// A bundle compiled to a binary has no entrypoints: the binary is what it is, and what
// the runtime starts to serve it (novox/hq ADR 0193). So a Go tools bundle is served
// as Go — the runtime execs it — exactly as a TypeScript one is through its launcher.
if bin := BinaryOf(a); bin != "" {
loads = []string{bin}
}
} }
// **Kept, never routed** (ADR 0155): the builder publishes to the store at the address // **Kept, never routed** (ADR 0155): the builder publishes to the store at the address
// it reached it by, and a manifest carrying that address names an installation — // it reached it by, and a manifest carrying that address names an installation —
@@ -220,11 +208,6 @@ func (b *Build) problems(module string) []string {
// A bundle's source is the module's own directory by definition, and what it needs to say // A bundle's source is the module's own directory by definition, and what it needs to say
// is which compiler — because the mesh chooses that, and cannot choose for a module that // is which compiler — because the mesh chooses that, and cannot choose for a module that
// has not said. // has not said.
if len(a.External) > 0 && (a.Kind != ArtifactBundle || a.Language != "typescript") {
problems = append(problems, fmt.Sprintf(
"%s: %q names packages it keeps external, and only a TypeScript bundle is bundled into "+
"one file with some kept out (novox/hq ADR 0193)", module, a.Name))
}
if len(a.Env) > 0 && a.Kind != ArtifactBundle { if len(a.Env) > 0 && a.Kind != ArtifactBundle {
problems = append(problems, fmt.Sprintf( problems = append(problems, fmt.Sprintf(
"%s: %q is a %q and says what it is given (env). Only a bundle the node's runtime "+ "%s: %q is a %q and says what it is given (env). Only a bundle the node's runtime "+
@@ -259,11 +242,6 @@ func (b *Build) problems(module string) []string {
for _, e := range a.Entrypoints { for _, e := range a.Entrypoints {
found = found || e == load found = found || e == load
} }
// A bundle compiled to a binary is one executable: the runtime loads that or nothing
// (novox/hq ADR 0193).
if bin := BinaryOf(a); bin != "" {
found = load == bin
}
if !found { if !found {
problems = append(problems, fmt.Sprintf( problems = append(problems, fmt.Sprintf(
"%s: %q says the runtime loads %q, which is not among its entrypoints — "+ "%s: %q says the runtime loads %q, which is not among its entrypoints — "+
@@ -458,40 +436,3 @@ func BinaryOf(a Artifact) string {
} }
return a.Name return a.Name
} }
// runByAResource is the artifacts one of a module's own resources names — a process that runs it,
// a step, an archive that unpacks it — by name.
func runByAResource(m Manifest) map[string]bool {
named := map[string]bool{}
for _, r := range m.Resources {
if a, ok := r["artifact"].(string); ok && a != "" {
named[a] = true
}
}
return named
}
// undeliveredBundles says which of a module's bundles nothing would ever put on a machine (novox/hq
// 04-ISSUES/216). A bundle reaches a machine three ways: the node's runtime serves it (it says
// `loads`, or its module declares `tools`), a resource names it (a process, a step, an archive), or
// it is the runtime itself. One reached by none of them was built, recorded and pushed as success,
// and was simply absent — seven modules' tools went missing that way on 2026-10-03. Refused here,
// naming the field that would deliver it.
func undeliveredBundles(m Manifest) []string {
if m.Build == nil || m.Module == RuntimeModule {
return nil
}
named := runByAResource(m)
var problems []string
for _, a := range m.Build.Artifacts {
if a.Kind != ArtifactBundle || named[a.Name] || len(a.Loads) > 0 || len(m.Tools) > 0 {
continue
}
problems = append(problems, fmt.Sprintf(
"%s: the bundle %q would be built and never reach a machine: nothing loads it, runs it or "+
"unpacks it. A tools bundle says `loads` (the entrypoints the node's runtime serves) or its "+
"module lists its `tools`; a daemon or a step is a resource naming it (novox/hq 04-ISSUES/216)",
m.Module, a.Name))
}
return problems
}
+4 -52
View File
@@ -693,15 +693,7 @@ func (r Resolution) compose(with Rendering, owner map[string]string,
// Now, and not before: a module whose resources are computed replaces them wholesale, and // Now, and not before: a module whose resources are computed replaces them wholesale, and
// merging earlier would throw away the files it still needs. // merging earlier would throw away the files it still needs.
// resources = append(append([]map[string]any{}, first...), resources...)
// **Except the module's own accounts, which go before even those** (novox/hq issue 213). What
// the mesh computes may belong to one: a module whose code runs as an account it declares has
// its secrets written owned by that account, and a file given to a user the machine does not
// have yet fails — so on the first apply the secrets were refused, the process started without
// them, and the second apply healed it, which is the fault the paragraph above describes.
// An account depends on nothing the mesh computes.
accounts, rest := accountsFirst(resources)
resources = append(append(accounts, first...), rest...)
// No container is given the mesh's names (novox/hq ADR 0148). It used to be: every // No container is given the mesh's names (novox/hq ADR 0148). It used to be: every
// container got the whole roster as `--add-host` entries at creation, and a name that // container got the whole roster as `--add-host` entries at creation, and a name that
@@ -880,12 +872,6 @@ func (r Resolution) compose(with Rendering, owner map[string]string,
if renamed := reflectsRenamed(m.Module, resource["reload-on"]); renamed != nil { if renamed := reflectsRenamed(m.Module, resource["reload-on"]); renamed != nil {
copied["reload-on"] = renamed copied["reload-on"] = renamed
} }
// And what a process replaces (novox/hq issue 213): a resource of this module's that it
// no longer declares, named as the host recorded it, or the host hands nothing over and
// removes it first.
if renamed := reflectsRenamed(m.Module, resource["replaces"]); renamed != nil {
copied["replaces"] = renamed
}
// **What reads one of this module's own secrets is restarted when it changes** (novox/hq // **What reads one of this module's own secrets is restarted when it changes** (novox/hq
// issue 203, issue 206). A credential is re-issued by the mesh, and a container that // issue 203, issue 206). A credential is re-issued by the mesh, and a container that
// mounted the old file keeps the old one open: the build machine ran for an hour on a // mounted the old file keeps the old one open: the build machine ran for an hour on a
@@ -2024,18 +2010,12 @@ func preparationTarget(m Manifest) string {
return "" return ""
} }
for _, r := range m.Resources { for _, r := range m.Resources {
// A container, or a process the host runs from a bundle the module built (novox/hq issue if fmt.Sprint(r["type"]) != "container" || !ownArtifact(r, m.Module) {
// 213): the same program in the same context, hosted as a unit rather than a container.
kind := fmt.Sprint(r["type"])
if (kind != "container" && kind != "process") || !ownArtifact(r, m.Module) {
continue continue
} }
if once, _ := r["run-once"].(bool); once { if once, _ := r["run-once"].(bool); once {
continue continue
} }
if r["schedule"] != nil {
continue
}
return fmt.Sprint(r["id"]) return fmt.Sprint(r["id"])
} }
return "" return ""
@@ -2049,10 +2029,7 @@ func ownArtifact(resource map[string]any, module string) bool {
return true return true
} }
image, _ := resource["image"].(string) image, _ := resource["image"].(string)
// A process or an archive carries what was built as its source (novox/hq issue 213). return strings.HasPrefix(image, ArtifactStoreScheme+module+"/")
source, _ := resource["source"].(string)
return strings.HasPrefix(image, ArtifactStoreScheme+module+"/") ||
strings.HasPrefix(source, ArtifactStoreScheme+module+"/")
} }
// prepared is the module's own resource as the step that prepares its state: the same image, the same // prepared is the module's own resource as the step that prepares its state: the same image, the same
@@ -2074,19 +2051,7 @@ func prepared(from map[string]any) map[string]any {
step["id"] = fmt.Sprint(from["id"]) + "-prepare" step["id"] = fmt.Sprint(from["id"]) + "-prepare"
step["name"] = fmt.Sprint(from["name"]) + "-prepare" step["name"] = fmt.Sprint(from["name"]) + "-prepare"
step["run-once"] = true step["run-once"] = true
if fmt.Sprint(from["type"]) == "process" { step["args"] = []any{PreparationArgument}
// A process says its whole command: the program, then its arguments. The step is the same
// program asked to prepare (novox/hq issue 213). It replaces nothing — what the process
// replaces is handed over to the process, never to the step that runs before it — and a
// step is not restarted, it runs again when what it reads changed, which `restart-on` says.
run := stringsIn(from["run"])
if len(run) > 0 {
step["run"] = []any{run[0], PreparationArgument}
}
delete(step, "replaces")
} else {
step["args"] = []any{PreparationArgument}
}
delete(step, "ports") delete(step, "ports")
delete(step, "ip") delete(step, "ip")
delete(step, "schedule") delete(step, "schedule")
@@ -2231,16 +2196,3 @@ func withRestartOn(have any, add []string) []any {
} }
return out return out
} }
// accountsFirst splits a module's resources into its accounts and everything else, each in the order
// written.
func accountsFirst(resources []map[string]any) (accounts, rest []map[string]any) {
for _, r := range resources {
if fmt.Sprint(r["type"]) == "user" {
accounts = append(accounts, r)
continue
}
rest = append(rest, r)
}
return accounts, rest
}
+17 -23
View File
@@ -1,7 +1,6 @@
package catalogue package catalogue
import ( import (
"encoding/json"
"fmt" "fmt"
"os" "os"
"reflect" "reflect"
@@ -171,7 +170,7 @@ func TestTheForgeHoldsTheNpmAndGitSeats(t *testing.T) {
// **And the forge's own address follows it**, composed from the manifest in the catalogue beside // **And the forge's own address follows it**, composed from the manifest in the catalogue beside
// this checkout (novox/hq 04-ISSUES/088). // this checkout (novox/hq 04-ISSUES/088).
// //
// The forge is reached a third way that neither test above covers: by its own code, over the // The forge is reached a third way that neither test above covers: by its own sidecar, over the
// machine's loopback, told where to go in its environment. The `2999:3000` mapping that lets the // machine's loopback, told where to go in its environment. The `2999:3000` mapping that lets the
// forge go on binding 3000 does nothing for a caller dialling the machine — so a literal there is // forge go on binding 3000 does nothing for a caller dialling the machine — so a literal there is
// wrong on every node whose assignment differs, and wrong for a second reason on a node given the // wrong on every node whose assignment differs, and wrong for a second reason on a node given the
@@ -179,13 +178,13 @@ func TestTheForgeHoldsTheNpmAndGitSeats(t *testing.T) {
// in an `env` at all is a declaration, not a manifest. // in an `env` at all is a declaration, not a manifest.
func TestTheForgesOwnAddressFollowsThePortTheNodeGaveIt(t *testing.T) { func TestTheForgesOwnAddressFollowsThePortTheNodeGaveIt(t *testing.T) {
forge, err := catalogueManifest(t, "gitea").Resolve([]Built{{ forge, err := catalogueManifest(t, "gitea").Resolve([]Built{{
Name: "code", Kind: ArtifactBundle, Name: "runtime", Kind: ArtifactImage,
Reference: ArtifactStoreScheme + "gitea/code/blobs/" + bundleDigest, Digest: bundleDigest, Reference: "registry.example/gitea-runtime@sha256:" + strings.Repeat("a", 64),
}}) }})
if err != nil { if err != nil {
t.Fatalf("the forge's manifest does not resolve against its own build: %v", err) t.Fatalf("the forge's manifest does not resolve against its own build: %v", err)
} }
r := Resolution{Node: "anchor", Modules: []Manifest{forge, theRuntime(t)}, Needs: []Needed{ r := Resolution{Node: "anchor", Modules: []Manifest{forge}, Needs: []Needed{
{Name: "postgres-database", For: "gitea", From: "anchor", At: "127.0.0.1", {Name: "postgres-database", For: "gitea", From: "anchor", At: "127.0.0.1",
Serves: map[string]any{"port": float64(5432)}, Sealed: "sealed-db"}, Serves: map[string]any{"port": float64(5432)}, Sealed: "sealed-db"},
{Name: "route", For: "gitea", From: "anchor"}, {Name: "route", For: "gitea", From: "anchor"},
@@ -195,8 +194,8 @@ func TestTheForgesOwnAddressFollowsThePortTheNodeGaveIt(t *testing.T) {
// The number this node was given for the forge — the one the machine it is about to run on // The number this node was given for the forge — the one the machine it is about to run on
// already publishes. // already publishes.
out, err := r.Declaration(Rendering{ArtifactStore: "anchor.internal:5101", out, err := r.Declaration(Rendering{
Needed: map[string]map[string]string{RuntimeModule: {"broker": "sealed-broker"}}, Needed: map[string]map[string]string{"gitea": {"broker": "sealed-broker"}},
Given: map[string]map[int]int{"gitea": {3000: 2999}}, Given: map[string]map[int]int{"gitea": {3000: 2999}},
}) })
if err != nil { if err != nil {
@@ -211,19 +210,14 @@ func TestTheForgesOwnAddressFollowsThePortTheNodeGaveIt(t *testing.T) {
if published := fmt.Sprint(server["ports"]); !strings.Contains(published, "2999:3000") { if published := fmt.Sprint(server["ports"]); !strings.Contains(published, "2999:3000") {
t.Fatalf("the forge is not published on the port this node gave it: %v", server["ports"]) t.Fatalf("the forge is not published on the port this node gave it: %v", server["ports"])
} }
// The forge's own code runs in the node's runtime (novox/hq ADR 0198), given its words there. runtime := fileNamed(out, "gitea.runtime")
runtime := fileNamed(out, RuntimeModule+"."+RuntimeProcessID())
if runtime == nil { if runtime == nil {
t.Fatalf("the node's runtime is not in the declaration: %v", ids(out)) t.Fatalf("the forge's sidecar is not in the declaration: %v", out)
} }
env, _ := runtime["env"].(map[string]string) env, _ := runtime["env"].(map[string]any)
var given map[string]map[string]string if env["MESH_GITEA_URL"] != "http://127.0.0.1:2999" {
if err := json.Unmarshal([]byte(env[RuntimeToolEnv]), &given); err != nil { t.Fatalf("the forge's sidecar dials %v while the machine publishes the forge on 2999 — "+
t.Fatalf("the runtime's %s is not JSON: %q", RuntimeToolEnv, env[RuntimeToolEnv]) "whatever reads it dials a dead port", env["MESH_GITEA_URL"])
}
if given["gitea"]["MESH_GITEA_URL"] != "http://127.0.0.1:2999" {
t.Fatalf("the forge's code dials %v while the machine publishes the forge on 2999 — "+
"whatever reads it dials a dead port", given["gitea"]["MESH_GITEA_URL"])
} }
} }
@@ -235,13 +229,13 @@ func declaredGiteaSsh(t *testing.T, given map[int]int) map[string]any {
t.Helper() t.Helper()
forge := catalogueManifest(t, "gitea") forge := catalogueManifest(t, "gitea")
resolved, err := forge.Resolve([]Built{{ resolved, err := forge.Resolve([]Built{{
Name: "code", Kind: ArtifactBundle, Name: "runtime", Kind: ArtifactImage,
Reference: ArtifactStoreScheme + "gitea/code/blobs/" + bundleDigest, Digest: bundleDigest, Reference: "registry.example/gitea-runtime@sha256:" + strings.Repeat("a", 64),
}}) }})
if err != nil { if err != nil {
t.Fatalf("the forge's manifest does not resolve against its own build: %v", err) t.Fatalf("the forge's manifest does not resolve against its own build: %v", err)
} }
r := Resolution{Node: "anchor", Modules: []Manifest{resolved, theRuntime(t)}, Needs: []Needed{ r := Resolution{Node: "anchor", Modules: []Manifest{resolved}, Needs: []Needed{
{Name: "postgres-database", For: "gitea", From: "anchor", At: "127.0.0.1", {Name: "postgres-database", For: "gitea", From: "anchor", At: "127.0.0.1",
Serves: map[string]any{"port": float64(5432)}, Sealed: "sealed-db"}, Serves: map[string]any{"port": float64(5432)}, Sealed: "sealed-db"},
{Name: "route", For: "gitea", From: "anchor"}, {Name: "route", For: "gitea", From: "anchor"},
@@ -252,8 +246,8 @@ func declaredGiteaSsh(t *testing.T, given map[int]int) map[string]any {
for k, v := range given { for k, v := range given {
givenPorts[k] = v givenPorts[k] = v
} }
out, err := r.Declaration(Rendering{ArtifactStore: "anchor.internal:5101", out, err := r.Declaration(Rendering{
Needed: map[string]map[string]string{RuntimeModule: {"broker": "sealed-broker"}}, Needed: map[string]map[string]string{"gitea": {"broker": "sealed-broker"}},
Ports: map[string]map[int]int{"gitea": givenPorts}, Ports: map[string]map[int]int{"gitea": givenPorts},
Given: map[string]map[int]int{"gitea": given}, Given: map[string]map[int]int{"gitea": given},
}) })
+1 -54
View File
@@ -757,11 +757,6 @@ type Artifact struct {
// list twice. A module declaring no tools has nothing the runtime loads, whatever it compiles. // list twice. A module declaring no tools has nothing the runtime loads, whatever it compiles.
Loads []string `json:"loads,omitempty"` Loads []string `json:"loads,omitempty"`
// External are packages a TypeScript bundle keeps as imports rather than inlining — a native
// addon, a package that reads its own files — and so carries the toolchain's node_modules for
// (novox/hq ADR 0193). Absent for nearly every bundle, which is then one file per entrypoint.
External []string `json:"external,omitempty"`
// Env is what a tools bundle is given on a machine (novox/hq ADR 0192): words and their values, // Env is what a tools bundle is given on a machine (novox/hq ADR 0192): words and their values,
// paths and constants composed with ${dir:…} and ${port:…} exactly as a container's environment // paths and constants composed with ${dir:…} and ${port:…} exactly as a container's environment
// is, never a secret's content. The node's runtime hands it to this bundle and to no other. // is, never a secret's content. The node's runtime hands it to this bundle and to no other.
@@ -1378,7 +1373,6 @@ func ParseManifest(raw []byte) (Manifest, error) {
} }
} }
problems = append(problems, m.Build.problems(m.Module)...) problems = append(problems, m.Build.problems(m.Module)...)
problems = append(problems, undeliveredBundles(m)...)
// **What provides the artifact store cannot be delivered through it** (novox/hq 04-ISSUES/029). // **What provides the artifact store cannot be delivered through it** (novox/hq 04-ISSUES/029).
// //
// Building publishes to the store, and the builder will not start without one. So a module // Building publishes to the store, and the builder will not start without one. So a module
@@ -1515,53 +1509,6 @@ func ParseManifest(raw []byte) (Manifest, error) {
"program that reads what the mesh delivered and reconciles", "program that reads what the mesh delivered and reconciles",
m.Module, r["id"])) m.Module, r["id"]))
} }
// **What a process replaces is something the module no longer declares** (novox/hq issue 213).
// The host keeps it running until the process is, then removes it: so it is named by the id the
// module used to give it, it is never a resource the module still declares — that would be
// applied and removed by one declaration — and only a process that stays up has anything to
// hand over to. Said here, near the author, as the host would refuse it far away.
ids := map[string]bool{}
for _, r := range m.Resources {
ids[fmt.Sprint(r["id"])] = true
}
for _, r := range m.Resources {
raw, present := r["replaces"]
if !present {
continue
}
if fmt.Sprint(r["type"]) != "process" {
problems = append(problems, fmt.Sprintf(
"%s: %v says what it replaces, and only a process does", m.Module, r["id"]))
continue
}
if once, _ := r["run-once"].(bool); once || r["schedule"] != nil {
problems = append(problems, fmt.Sprintf(
"%s: %v replaces something and runs once or on a schedule — only a process that stays "+
"up is there a moment later to hand over to", m.Module, r["id"]))
}
list, ok := raw.([]any)
if !ok {
problems = append(problems, fmt.Sprintf(
"%s: %v replaces %v; replaces is a list of the ids this module no longer declares",
m.Module, r["id"], raw))
continue
}
for _, item := range list {
id, ok := item.(string)
switch {
case !ok || strings.TrimSpace(id) == "":
problems = append(problems, fmt.Sprintf(
"%s: %v replaces %v, which is not an id", m.Module, r["id"], item))
case strings.Contains(id, "."):
problems = append(problems, fmt.Sprintf(
"%s: %v replaces %q; a process replaces only a resource of its own module, named "+
"by its own id", m.Module, r["id"], id))
case ids[id]:
problems = append(problems, fmt.Sprintf(
"%s: %v replaces %q, which this module still declares", m.Module, r["id"], id))
}
}
}
// **A module that prepares its state must have code the mesh can run** (novox/hq ADR 0135). The // **A module that prepares its state must have code the mesh can run** (novox/hq ADR 0135). The
// preparation is the module's own program in its preparation mode, so it is derived from the // preparation is the module's own program in its preparation mode, so it is derived from the
// resource that runs that program — and a module declaring none has asked for something the mesh // resource that runs that program — and a module declaring none has asked for something the mesh
@@ -1569,7 +1516,7 @@ func ParseManifest(raw []byte) (Manifest, error) {
// quietly prepares nothing. // quietly prepares nothing.
if m.Prepares && preparationTarget(m) == "" { if m.Prepares && preparationTarget(m) == "" {
problems = append(problems, fmt.Sprintf( problems = append(problems, fmt.Sprintf(
"%s says it prepares its state, and declares no container or process running an artifact it built — "+ "%s says it prepares its state, and declares no container running an artifact it built — "+
"the preparation is this module's own program, so there has to be one for the mesh to "+ "the preparation is this module's own program, so there has to be one for the mesh to "+
"run it in", m.Module)) "run it in", m.Module))
} }
+5 -10
View File
@@ -31,7 +31,7 @@ import (
// //
// So a module asks. `${port:8080}` is "the machine-side port you gave me for the 8080 I said I // So a module asks. `${port:8080}` is "the machine-side port you gave me for the 8080 I said I
// listen on", and the module writes that where it would otherwise have written a literal — in a // listen on", and the module writes that where it would otherwise have written a literal — in a
// file's content, or in a value of a container's or a process's `env`. // file's content, or in a value of a container's `env`.
// //
// **The environment is filled by the control plane, exactly as a bound value is.** A port is not // **The environment is filled by the control plane, exactly as a bound value is.** A port is not
// secret — the mesh holds it in the clear — so there is nothing for the host to be the only // secret — the mesh holds it in the clear — so there is nothing for the host to be the only
@@ -64,12 +64,7 @@ func portsUsed(content string) []int {
} }
// portInto replaces a resource's ${port:…} placeholders with what this machine assigned — in a // portInto replaces a resource's ${port:…} placeholders with what this machine assigned — in a
// file's content, and in a value of a container's or a process's environment. // file's content, and in a value of a container's environment.
//
// **A process's environment is a container's** (novox/hq to-be 38 WP4c). A module's code moving out
// of its container becomes a process on the machine and still has to be told what the container
// was told; filled for one kind and not the other, the literal reached the process and was read as
// a port, and the modules that moved first wrote their run-once steps a 0600 env file instead.
// //
// A port the module did not say it listens on is refused, for the same reason a binding's unknown // A port the module did not say it listens on is refused, for the same reason a binding's unknown
// key is: the module is asking about something it never declared, and the answer would be a guess. // key is: the module is asking about something it never declared, and the answer would be a guess.
@@ -89,7 +84,7 @@ func portInto(resource map[string]any, module string, listens []Listening, with
} }
resource["content"] = filled resource["content"] = filled
case "container", "process": case "container":
env, ok := resource["env"].(map[string]any) env, ok := resource["env"].(map[string]any)
if !ok { if !ok {
return nil return nil
@@ -111,8 +106,8 @@ func portInto(resource map[string]any, module string, listens []Listening, with
continue continue
} }
value, err := portsFilledInto(written, value, err := portsFilledInto(written,
fmt.Sprintf("%s's %s %s sets %s to something that", fmt.Sprintf("%s's container %s sets %s to something that",
module, resource["type"], resource["name"], key), module, listens, with) module, resource["name"], key), module, listens, with)
if err != nil { if err != nil {
return err return err
} }
-80
View File
@@ -1,80 +0,0 @@
package catalogue
import (
"strings"
"testing"
)
// **A process's environment is composed as a container's is** (novox/hq to-be 38 WP4c).
//
// A module's code moving out of its container becomes a process on the machine, and what its
// container's environment asked for — the port this machine gave the module, the place it put the
// module's directory — it still has to be told. Filled for a container and not for a process, the
// literal `${port:8080}` reached the process as its environment and was read as a port; the modules
// that moved first wrote their run-once steps an env file instead.
func processModule(env map[string]any) Manifest {
return Manifest{
Module: "showcase",
Listens: []Listening{{Port: 8080, From: FromMesh}},
Resources: []map[string]any{
{"id": "data", "type": "directory", "mode": "0700"},
{"id": "setup", "type": "process", "name": "showcase-setup", "run-once": true,
"run": []any{"/usr/bin/showcase", "setup"}, "env": env},
},
}
}
func TestAProcessIsToldItsPortAndItsPlaceInItsEnvironment(t *testing.T) {
env := map[string]any{
"SHOWCASE_URL": "http://127.0.0.1:${port:8080}",
"SHOWCASE_DATA": "${dir:data}/objects",
"SHOWCASE_DB": "127.0.0.1:${seat:mesh-store:5432}",
"GREETING": "hello",
}
out, err := Resolution{Node: "anchor", Modules: []Manifest{processModule(env)}}.Declaration(Rendering{
Ports: map[string]map[int]int{"showcase": {8080: 21000}},
Seats: map[string]map[int]int{"mesh-store": {5432: 6852}},
})
if err != nil {
t.Fatalf("a process asking for its port and its place does not compose: %v", err)
}
setup := fileNamed(out, "showcase.setup")
if setup == nil {
t.Fatalf("the process is not in the declaration: %v", out)
}
got, _ := setup["env"].(map[string]any)
for key, want := range map[string]string{
"SHOWCASE_URL": "http://127.0.0.1:21000",
"SHOWCASE_DATA": "/var/lib/showcase/data/objects",
"SHOWCASE_DB": "127.0.0.1:6852",
"GREETING": "hello",
} {
if got[key] != want {
t.Errorf("the process is told %s=%v, want %q", key, got[key], want)
}
}
if env["SHOWCASE_URL"] != "http://127.0.0.1:${port:8080}" {
t.Fatalf("composing for one machine edited the module's own manifest: %v", env)
}
}
// An unknown reference in a process's environment is refused as a container's is, naming the
// process and the variable — left alone, it would reach the machine as a literal.
func TestAProcessAskingAboutAnUndeclaredPortIsRefused(t *testing.T) {
env := map[string]any{"SHOWCASE_URL": "http://127.0.0.1:${port:9999}"}
_, err := Resolution{Node: "anchor", Modules: []Manifest{processModule(env)}}.Declaration(Rendering{})
if err == nil {
t.Fatal("a process was told a port its module never said it listens on")
}
for _, said := range []string{"showcase-setup", "SHOWCASE_URL", "${port:9999}", "8080"} {
if !strings.Contains(err.Error(), said) {
t.Errorf("the refusal does not say %q: %v", said, err)
}
}
env = map[string]any{"SHOWCASE_DATA": "${dir:date}/objects"}
if _, err := (Resolution{Node: "anchor", Modules: []Manifest{processModule(env)}}).Declaration(Rendering{}); err == nil ||
!strings.Contains(err.Error(), "${dir:date}") {
t.Fatalf("a process naming no directory of its module was not refused: %v", err)
}
}
-68
View File
@@ -389,71 +389,3 @@ func TestARuntimeCompiledToABinaryRunsItself(t *testing.T) {
t.Errorf("the Go runtime is not told what to serve or whose it is: %v %v", env, process["user"]) t.Errorf("the Go runtime is not told what to serve or whose it is: %v %v", env, process["user"])
} }
} }
// novox/hq 04-ISSUES/216: a bundle nothing loads, runs or unpacks is refused at registration; saying
// `loads`, listing `tools`, or a resource naming it admits it.
func TestABundleNothingDeliversIsRefused(t *testing.T) {
base := func() Manifest {
return Manifest{Module: "baserow", Version: "1", Build: &Build{Artifacts: []Artifact{
{Name: "tools", Kind: ArtifactBundle, Language: "typescript", Entrypoints: []string{"tools/index.js"}}}}}
}
if p := undeliveredBundles(base()); len(p) != 1 || !strings.Contains(p[0], "never reach a machine") {
t.Fatalf("a bundle nothing delivers was admitted: %v", p)
}
loads := base()
loads.Build.Artifacts[0].Loads = []string{"tools/index.js"}
tools := base()
tools.Tools = []string{"baserow_list_rows"}
run := base()
run.Resources = []map[string]any{{"id": "daemon", "type": "process", "artifact": "tools", "run": []any{"node", "tools/index.js"}}}
runtime := base()
runtime.Module = RuntimeModule
for name, m := range map[string]Manifest{"loads": loads, "tools": tools, "a process": run, "the runtime": runtime} {
if p := undeliveredBundles(m); len(p) != 0 {
t.Errorf("a bundle delivered by %s was refused: %v", name, p)
}
}
}
// novox/hq ADR 0193: a Go tools bundle is served — its binary is what the runtime starts, delivered
// like any tools bundle, named to the runtime where a TypeScript bundle names its launcher.
func TestAGoToolsBundleIsServedByItsBinary(t *testing.T) {
with := Rendering{ArtifactStore: "anchor.internal:5101",
Needed: map[string]map[string]string{RuntimeModule: {"broker": "sealed-credential"}}}
lamp := Manifest{Module: "lamp", Version: "1", Tools: []string{"on"},
Build: &Build{Artifacts: []Artifact{{Name: "tools", Kind: ArtifactBundle, Language: "go",
System: "arch", From: "cmd/lamp-tools"}}}}
if p := lamp.Build.problems("lamp"); len(p) != 0 {
t.Fatalf("a Go tools bundle was refused: %v", p)
}
lamp, err := lamp.Resolve([]Built{{Name: "tools", Kind: ArtifactBundle,
Reference: ArtifactStoreScheme + "lamp/tools/blobs/" + bundleDigest, Digest: bundleDigest}})
if err != nil {
t.Fatal(err)
}
if fmt.Sprint(lamp.Bundles[0].Loads) != "[lamp-tools]" {
t.Fatalf("the runtime loads %v from a Go bundle, want its binary", lamp.Bundles[0].Loads)
}
out, err := Resolution{Node: "anchor", Account: "ops", Modules: []Manifest{lamp, theRuntime(t)}}.Declaration(with)
if err != nil {
t.Fatal(err)
}
if fileNamed(out, "lamp."+BundleID("tools")) == nil {
t.Errorf("the Go bundle is not delivered: %v", ids(out))
}
env := fileNamed(out, RuntimeModule+"."+RuntimeProcessID())["env"].(map[string]string)
if env[RuntimeToolModules] != "lamp="+BundlePath("lamp", "tools")+"/lamp-tools" {
t.Errorf("the runtime is told %q, want the binary", env[RuntimeToolModules])
}
// An artifact may say it explicitly; naming anything but the binary is refused.
said := Manifest{Module: "lamp", Version: "1", Build: &Build{Artifacts: []Artifact{{Name: "tools",
Kind: ArtifactBundle, Language: "go", System: "arch", Binary: "lamp", Loads: []string{"lamp"}}}}}
if p := said.Build.problems("lamp"); len(p) != 0 {
t.Errorf("loads naming the binary was refused: %v", p)
}
said.Build.Artifacts[0].Loads = []string{"tools/index.js"}
if p := said.Build.problems("lamp"); len(p) == 0 {
t.Error("a Go bundle loading a file it does not contain was admitted")
}
}
+5 -5
View File
@@ -43,8 +43,8 @@ import (
var ofSeat = regexp.MustCompile(`\$\{seat:([a-z0-9][a-z0-9-]*):([0-9]+)\}`) var ofSeat = regexp.MustCompile(`\$\{seat:([a-z0-9][a-z0-9-]*):([0-9]+)\}`)
// seatInto replaces a resource's ${seat:…} placeholders with where this machine put each seat's // seatInto replaces a resource's ${seat:…} placeholders with where this machine put each seat's
// holder — in a file's content, and in a value of a container's or a process's environment. The // holder — in a file's content, and in a value of a container's environment. The same two places
// same places portInto fills, for the same reason: they are where a program reads a number from. // portInto fills, for the same reason: they are where a process reads a number from.
func seatInto(resource map[string]any, module string, with Rendering) error { func seatInto(resource map[string]any, module string, with Rendering) error {
switch fmt.Sprint(resource["type"]) { switch fmt.Sprint(resource["type"]) {
case "file": case "file":
@@ -58,7 +58,7 @@ func seatInto(resource map[string]any, module string, with Rendering) error {
} }
resource["content"] = filled resource["content"] = filled
case "container", "process": case "container":
env, ok := resource["env"].(map[string]any) env, ok := resource["env"].(map[string]any)
if !ok { if !ok {
return nil return nil
@@ -78,8 +78,8 @@ func seatInto(resource map[string]any, module string, with Rendering) error {
continue continue
} }
value, err := seatsFilledInto(written, value, err := seatsFilledInto(written,
fmt.Sprintf("%s's %s %s sets %s to something that", fmt.Sprintf("%s's container %s sets %s to something that",
module, resource["type"], resource["name"], key), with) module, resource["name"], key), with)
if err != nil { if err != nil {
return err return err
} }
+14 -43
View File
@@ -42,29 +42,9 @@ type Build struct {
// Failed is the builder's own words, empty when it worked. // Failed is the builder's own words, empty when it worked.
Failed string Failed string
Made []Artifact Made []Artifact
// Asked is when the build was requested, zero when that is not known (an id of another shape, At time.Time
// or a build recorded before the mesh kept it). **What orders one build of a module against
// another** (novox/hq 04-ISSUES/219): builds in flight together finish in any order, and the
// one asked last stood on the newest bases.
Asked time.Time
// At is when the outcome was recorded — when it finished, not when it was asked.
At time.Time
} }
// AskedOrAt is when the build was asked, or when it was recorded when that is not known — the
// order the mesh had before it kept the request time.
func (b Build) AskedOrAt() time.Time {
if !b.Asked.IsZero() {
return b.Asked
}
return b.At
}
// newestRequestFirst is the ordering every "what a module currently is" question uses: the newest
// request wins, whenever it finished (novox/hq 04-ISSUES/219). A build whose request time is not
// known is placed at the moment it was recorded, which is the rule that held before.
const newestRequestFirst = `coalesce(asked, at) desc, at desc`
// ReadRepository is a repository a build read source from besides the module's own. // ReadRepository is a repository a build read source from besides the module's own.
type ReadRepository struct { type ReadRepository struct {
Repository string `json:"repository"` Repository string `json:"repository"`
@@ -103,17 +83,13 @@ func (i *Inventory) RecordBuild(ctx context.Context, b Build) error {
if b.Module != "" { if b.Module != "" {
module = &b.Module module = &b.Module
} }
var asked *time.Time
if !b.Asked.IsZero() {
asked = &b.Asked
}
_, err = i.store.Pool().Exec(ctx, _, err = i.store.Pool().Exec(ctx,
`insert into build (id, repository, ref, module, commit_hash, built_on, failed, made, `insert into build (id, repository, ref, module, commit_hash, built_on, failed, made,
source_path, manifest, built_against, built_contexts, asked) source_path, manifest, built_against, built_contexts)
values ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13) values ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12)
on conflict (id) do nothing`, on conflict (id) do nothing`,
b.ID, b.Repository, b.Ref, module, b.Commit, b.On, b.Failed, made, b.ID, b.Repository, b.Ref, module, b.Commit, b.On, b.Failed, made,
b.Path, manifestOrNil(b.Manifest), against, read, asked) b.Path, manifestOrNil(b.Manifest), against, read)
return err return err
} }
@@ -126,11 +102,11 @@ func (i *Inventory) Builds(ctx context.Context, module string, limit int) ([]Bui
if limit <= 0 { if limit <= 0 {
limit = 20 limit = 20
} }
query := `select id, repository, ref, coalesce(module,''), commit_hash, built_on, failed, made, asked, at query := `select id, repository, ref, coalesce(module,''), commit_hash, built_on, failed, made, at
from build order by at desc limit $1` from build order by at desc limit $1`
args := []any{limit} args := []any{limit}
if module != "" { if module != "" {
query = `select id, repository, ref, coalesce(module,''), commit_hash, built_on, failed, made, asked, at query = `select id, repository, ref, coalesce(module,''), commit_hash, built_on, failed, made, at
from build where module = $2 order by at desc limit $1` from build where module = $2 order by at desc limit $1`
args = append(args, module) args = append(args, module)
} }
@@ -145,14 +121,10 @@ func (i *Inventory) Builds(ctx context.Context, module string, limit int) ([]Bui
for rows.Next() { for rows.Next() {
var b Build var b Build
var made []byte var made []byte
var asked *time.Time
if err := rows.Scan(&b.ID, &b.Repository, &b.Ref, &b.Module, &b.Commit, if err := rows.Scan(&b.ID, &b.Repository, &b.Ref, &b.Module, &b.Commit,
&b.On, &b.Failed, &made, &asked, &b.At); err != nil { &b.On, &b.Failed, &made, &b.At); err != nil {
return nil, err return nil, err
} }
if asked != nil {
b.Asked = *asked
}
if err := json.Unmarshal(made, &b.Made); err != nil { if err := json.Unmarshal(made, &b.Made); err != nil {
return nil, err return nil, err
} }
@@ -163,9 +135,8 @@ func (i *Inventory) Builds(ctx context.Context, module string, limit int) ([]Bui
// Held is every artifact this mesh has built, keyed "<module>/<artifact>". // Held is every artifact this mesh has built, keyed "<module>/<artifact>".
// //
// **The successful build of each module asked last wins**, which is the same rule the rest of the // **The newest successful build of each module wins**, which is the same rule the rest of the mesh
// mesh uses for what a module currently is — asked last, not finished last (novox/hq // uses for what a module currently is. A module rebuilt to something broken and then rebuilt again
// 04-ISSUES/219): an older request that finishes later stood on older bases. A module rebuilt to something broken and then rebuilt again
// is at the second one; a module whose last build failed is at the last one that worked, because a // is at the second one; a module whose last build failed is at the last one that worked, because a
// failure published nothing and the thing it published before is still what exists. // failure published nothing and the thing it published before is still what exists.
// //
@@ -176,7 +147,7 @@ func (i *Inventory) Held(ctx context.Context) (map[string]string, error) {
`select distinct on (module) module, made `select distinct on (module) module, made
from build from build
where module is not null and module <> '' and failed = '' where module is not null and module <> '' and failed = ''
order by module, `+newestRequestFirst) order by module, at desc`)
if err != nil { if err != nil {
return nil, err return nil, err
} }
@@ -214,7 +185,7 @@ func (i *Inventory) BuiltAgainst(ctx context.Context) (map[string][]string, erro
`select distinct on (module) module, built_against `select distinct on (module) module, built_against
from build from build
where module is not null and module <> '' and failed = '' where module is not null and module <> '' and failed = ''
order by module, `+newestRequestFirst) order by module, at desc`)
if err != nil { if err != nil {
return nil, err return nil, err
} }
@@ -252,7 +223,7 @@ func (i *Inventory) ReadRepositories(ctx context.Context) (map[string][]ReadRepo
`select distinct on (module) module, built_contexts `select distinct on (module) module, built_contexts
from build from build
where module is not null and module <> '' and failed = '' where module is not null and module <> '' and failed = ''
order by module, `+newestRequestFirst) order by module, at desc`)
if err != nil { if err != nil {
return nil, err return nil, err
} }
@@ -303,7 +274,7 @@ func manifestOrNil(raw []byte) any {
// follows when it decides whether to announce at all. // follows when it decides whether to announce at all.
// //
// One row per module and commit: a module built twice at the same commit is one fact, and the // One row per module and commit: a module built twice at the same commit is one fact, and the
// row asked last is the one whose artifacts are current (novox/hq 04-ISSUES/219). // latest row is the one whose artifacts are current.
func (i *Inventory) Announceable(ctx context.Context) ([]Build, error) { func (i *Inventory) Announceable(ctx context.Context) ([]Build, error) {
rows, err := i.store.Pool().Query(ctx, rows, err := i.store.Pool().Query(ctx,
`select distinct on (module, commit_hash) `select distinct on (module, commit_hash)
@@ -311,7 +282,7 @@ func (i *Inventory) Announceable(ctx context.Context) ([]Build, error) {
source_path, manifest, built_against, at source_path, manifest, built_against, at
from build from build
where failed = '' and module is not null and module <> '' and commit_hash <> '' where failed = '' and module is not null and module <> '' and commit_hash <> ''
order by module, commit_hash, `+newestRequestFirst) order by module, commit_hash, at desc`)
if err != nil { if err != nil {
return nil, err return nil, err
} }
+1 -39
View File
@@ -49,12 +49,6 @@ func (i *Inventory) BusRecords(ctx context.Context) (broker.Records, error) {
} }
} }
// Who holds each seat held once for the mesh, where the mesh recorded it (novox/hq issue 218).
holdings, err := i.Holdings(ctx)
if err != nil {
return broker.Records{}, fmt.Errorf("cannot read who holds the mesh's seats: %w", err)
}
out := broker.Records{Assigned: map[string][]broker.Declared{}, People: map[string][]string{}, out := broker.Records{Assigned: map[string][]broker.Declared{}, People: map[string][]string{},
Interchangeable: map[string]bool{}} Interchangeable: map[string]bool{}}
for _, n := range nodes { for _, n := range nodes {
@@ -78,9 +72,7 @@ func (i *Inventory) BusRecords(ctx context.Context) (broker.Records, error) {
"%s is assigned to %s and is not in the catalogue, so what it may say cannot "+ "%s is assigned to %s and is not in the catalogue, so what it may say cannot "+
"be derived", module, n.Name) "be derived", module, n.Name)
} }
d := declaredFor(m, seats) out.Assigned[n.Name] = append(out.Assigned[n.Name], declaredFor(m, seats))
d.Holds = heldHere(d.Holds, holdings, n.Name, module)
out.Assigned[n.Name] = append(out.Assigned[n.Name], d)
if m.Instances == catalogue.InstancesInterchangeable { if m.Instances == catalogue.InstancesInterchangeable {
out.Interchangeable[m.Module] = true out.Interchangeable[m.Module] = true
} }
@@ -191,33 +183,3 @@ func (i *Inventory) NodesWithALiveToken(ctx context.Context) ([]string, error) {
} }
return out, rows.Err() return out, rows.Err()
} }
// heldHere keeps of what a module claims only the seats it holds on this machine (novox/hq issue 218).
// A seat held once per machine is held by every assignment that claims it. A seat held once for the
// mesh is held by one assignment: where the mesh recorded who holds it, a claim on any other machine
// grants nothing and issues nothing — or the module would serve the role's verbs from a machine that
// is not the role's, and a question to the mesh's store would be answered from the wrong database. A
// mesh seat with no holder on record is left as it was derived.
func heldHere(claimed []broker.Seat, holdings []catalogue.Held, node, module string) []broker.Seat {
recorded := map[string][]catalogue.Held{}
for _, h := range holdings {
if h.Scope == catalogue.ScopeMesh {
recorded[h.Claim] = append(recorded[h.Claim], h)
}
}
var out []broker.Seat
for _, s := range claimed {
holders, onRecord := recorded[s.Name]
if s.Scope != catalogue.ScopeMesh || !onRecord {
out = append(out, s)
continue
}
for _, h := range holders {
if h.Node == node && h.Module == module {
out = append(out, s)
break
}
}
}
return out
}
+6 -38
View File
@@ -17,10 +17,6 @@ import (
// ErrNoSuchModule is what the mesh says about a module it has never been told about. // ErrNoSuchModule is what the mesh says about a module it has never been told about.
var ErrNoSuchModule = errors.New("no module of that name") var ErrNoSuchModule = errors.New("no module of that name")
// ErrSuperseded is a registration from a build asked before the one the module is already at
// (novox/hq 04-ISSUES/219). The build is recorded; what the module is does not change.
var ErrSuperseded = errors.New("a build asked later is already what the module is")
// ErrStillAssigned is why a module cannot be forgotten. // ErrStillAssigned is why a module cannot be forgotten.
// //
// Its own error because it is not a fault: it means a machine is running that module now, and // Its own error because it is not a fault: it means a machine is running that module now, and
@@ -51,10 +47,6 @@ type Source struct {
// itself no longer carries its build (novox/hq to-be 38 WP2.4). Empty for a manifest handed over // itself no longer carries its build (novox/hq to-be 38 WP2.4). Empty for a manifest handed over
// by hand, which carries its `build.on` itself. // by hand, which carries its `build.on` itself.
Against []string Against []string
// Asked is when the build this manifest came from was requested (novox/hq 04-ISSUES/219). Zero
// is a manifest handed over by hand, or a build whose request time is not known: either is
// taken as asked at the moment it is registered.
Asked time.Time
} }
// Current reports whether what the mesh holds is what the source last had. // Current reports whether what the mesh holds is what the source last had.
@@ -105,23 +97,13 @@ func (i *Inventory) RegisterModule(ctx context.Context, m catalogue.Manifest, fr
return err return err
} }
asked := from.Asked
if asked.IsZero() {
asked = time.Now()
}
// A module registered without provenance keeps whatever it had. Handing over a manifest by // A module registered without provenance keeps whatever it had. Handing over a manifest by
// hand is a legitimate way to fix something in a hurry, and it should not silently erase the // hand is a legitimate way to fix something in a hurry, and it should not silently erase the
// record of where the module normally comes from — which is the only thing that would say, // record of where the module normally comes from — which is the only thing that would say,
// afterwards, that the machine is running something nobody can rebuild. // afterwards, that the machine is running something nobody can rebuild.
// _, err = i.store.Pool().Exec(ctx,
// **An older request never replaces a newer one** (novox/hq 04-ISSUES/219). Builds of one `insert into module (name, manifest, version, source, source_path, source_seat, ref, built_from, source_head)
// module in flight together finish in any order, and each stood on the bases the mesh held when values ($1, $2, nullif($3,''), nullif($4,''), $7, $8, nullif($5,''), nullif($6,''), nullif($6,''))
// it was asked; the one asked later is what the module is, whichever is heard last. An outcome
// of an earlier request is kept in the build records and changes nothing here.
tag, err := i.store.Pool().Exec(ctx,
`insert into module (name, manifest, version, source, source_path, source_seat, ref, built_from, source_head, built_asked)
values ($1, $2, nullif($3,''), nullif($4,''), $7, $8, nullif($5,''), nullif($6,''), nullif($6,''), $9)
on conflict (name) do update set on conflict (name) do update set
manifest = excluded.manifest, manifest = excluded.manifest,
version = excluded.version, version = excluded.version,
@@ -133,23 +115,9 @@ func (i *Inventory) RegisterModule(ctx context.Context, m catalogue.Manifest, fr
else excluded.source_seat end, else excluded.source_seat end,
ref = coalesce(excluded.ref, module.ref), ref = coalesce(excluded.ref, module.ref),
built_from = coalesce(excluded.built_from, module.built_from), built_from = coalesce(excluded.built_from, module.built_from),
source_head = coalesce(excluded.built_from, module.source_head), source_head = coalesce(excluded.built_from, module.source_head)`,
built_asked = excluded.built_asked m.Module, raw, m.Version, from.Repository, from.Ref, from.BuiltFrom, from.Path, from.Seat)
where module.built_asked is null or module.built_asked <= excluded.built_asked`, return err
m.Module, raw, m.Version, from.Repository, from.Ref, from.BuiltFrom, from.Path, from.Seat, asked)
if err != nil {
return err
}
if tag.RowsAffected() == 0 {
var current time.Time
if err := i.store.Pool().QueryRow(ctx,
`select built_asked from module where name = $1`, m.Module).Scan(&current); err != nil {
return err
}
return fmt.Errorf("%w: %s is at a build asked %s, and this one was asked %s",
ErrSuperseded, m.Module, current.UTC().Format(time.RFC3339), asked.UTC().Format(time.RFC3339))
}
return nil
} }
// registeredInThatShape is whether the catalogue already holds this module as a tools container on // registeredInThatShape is whether the catalogue already holds this module as a tools container on
-20
View File
@@ -100,23 +100,3 @@ func TestABundleStandsOnTheToolchainItIsCompiledIn(t *testing.T) {
} }
} }
} }
// novox/hq 04-ISSUES/212: a toolchain standing on the SDK's package is planned after the SDK, so a
// release of the SDK rebuilds the toolchain, and every bundle compiled in it after that.
func TestAToolchainStandingOnTheSDKFollowsIt(t *testing.T) {
entries := []Entry{
{Manifest: catalogue.Manifest{Module: "mesh-sdk"}},
{Manifest: catalogue.Manifest{Module: "mesh-tools", Build: &catalogue.Build{
On: []catalogue.BuildsOn{{Arg: "MESH_SDK", Module: "mesh-sdk", Artifact: "lib"}}}}},
}
edges := dependenciesOf(entries, nil, nil)
found := false
for _, e := range edges {
if e.From == "mesh-tools" && e.To == "mesh-sdk" {
found = true
}
}
if !found {
t.Errorf("no edge from the toolchain to the SDK: %v", edges)
}
}
-32
View File
@@ -1,32 +0,0 @@
package inventory
import (
"testing"
"github.com/novox/mesh-controller/internal/broker"
"github.com/novox/mesh-controller/internal/catalogue"
)
// novox/hq issue 218: a seat held once for the mesh is granted and issued only to the holder on record;
// a node seat to every machine's claimant; a mesh seat with no holder on record as derived.
func TestOnlyTheRecordedHolderHoldsAMeshSeat(t *testing.T) {
claimed := []broker.Seat{
{Name: "mesh-store", Scope: catalogue.ScopeMesh},
{Name: "node-packet-filter", Scope: catalogue.ScopeNode},
{Name: "unrecorded", Scope: catalogue.ScopeMesh},
}
holdings := []catalogue.Held{{Claim: "mesh-store", Scope: catalogue.ScopeMesh, Node: "control", Module: "postgres"}}
names := func(ss []broker.Seat) (out []string) {
for _, s := range ss {
out = append(out, s.Name)
}
return
}
if got := names(heldHere(claimed, holdings, "control", "postgres")); len(got) != 3 {
t.Errorf("the holder lost a seat: %v", got)
}
got := names(heldHere(claimed, holdings, "other", "postgres"))
if len(got) != 2 || got[0] != "node-packet-filter" || got[1] != "unrecorded" {
t.Errorf("a claimant on another machine holds %v; want the node seat and the unrecorded one, not the store", got)
}
}
-59
View File
@@ -87,62 +87,3 @@ func (i *Inventory) tryHold(ctx context.Context, sorted []string) (func(), strin
} }
return release, "", nil return release, "", nil
} }
// ErrPlansBusy is the plans held by another act — on a machine replacing its controller, the other
// controller — for longer than a caller waits, or at all for one that does not wait.
var ErrPlansBusy = errors.New("another controller is working the plans")
// HoldPlans makes working the plans one act at a time, across every controller on the store
// (novox/hq issue 213). A plan is read, changed and written whole; two controllers doing that at
// once — the old and the new for the moment a machine hands its controller over, or a controller
// and a person's `plans stop` — each act on what the other has not saved yet: a tier asked twice,
// an outcome written over. A session-level advisory lock on one connection, released by the
// returned function and by the session ending, so a controller that dies holding it holds nothing.
//
// wait false gives ErrPlansBusy at once when another holds them — the timer's way: the holder is
// moving the plans already. wait true looks again every HoldPoll for up to HoldWaitFor — an
// outcome's or a merge's way, which must be written.
func (i *Inventory) HoldPlans(ctx context.Context, wait bool) (func(), error) {
deadline := time.Now().Add(HoldWaitFor)
for {
release, took, err := i.tryLock(ctx, "mesh-plans")
if err != nil || took {
return release, err
}
if !wait || time.Now().After(deadline) {
return nil, ErrPlansBusy
}
select {
case <-ctx.Done():
return nil, ctx.Err()
case <-time.After(HoldPoll):
}
}
}
// tryLock takes one named advisory lock on a connection of its own, or gives the connection back.
func (i *Inventory) tryLock(ctx context.Context, key string) (func(), bool, error) {
conn, err := i.store.Pool().Acquire(ctx)
if err != nil {
return nil, false, err
}
var once sync.Once
release := func() {
once.Do(func() {
if _, err := conn.Exec(context.WithoutCancel(ctx), `select pg_advisory_unlock_all()`); err != nil {
_ = conn.Conn().Close(context.WithoutCancel(ctx))
}
conn.Release()
})
}
var took bool
if err := conn.QueryRow(ctx, `select pg_try_advisory_lock(hashtext($1)::bigint)`, key).Scan(&took); err != nil {
release()
return nil, false, err
}
if !took {
release()
return nil, false, nil
}
return release, true, nil
}
-38
View File
@@ -1,38 +0,0 @@
package inventory
import (
"errors"
"testing"
"time"
)
// novox/hq issue 213: while a machine hands its controller over from the container to the process,
// two controllers run on one store for a moment. Working the plans is one act at a time across them.
func TestThePlansAreWorkedByOneControllerAtATime(t *testing.T) {
first := ForTest(t)
// A second controller: its own connections to the same store.
second, err := Open(t.Context())
if err != nil {
t.Fatal(err)
}
t.Cleanup(second.Close)
release, err := first.HoldPlans(t.Context(), false)
if err != nil {
t.Fatalf("the plans could not be held when nobody held them: %v", err)
}
t.Cleanup(release) // a pool closing waits for a connection still held; release is idempotent
if _, err := second.HoldPlans(t.Context(), false); !errors.Is(err, ErrPlansBusy) {
t.Fatalf("a second controller held the plans while the first did: %v", err)
}
// A waiter gets them once they are let go.
was := HoldPoll
HoldPoll = 10 * time.Millisecond
defer func() { HoldPoll = was }()
go func() { time.Sleep(50 * time.Millisecond); release() }()
again, err := second.HoldPlans(t.Context(), true)
if err != nil {
t.Fatalf("a waiting controller never got the plans once they were let go: %v", err)
}
again()
}
@@ -1,23 +0,0 @@
-- A build is ordered by when it was asked, not when it finished (novox/hq 04-ISSUES/219).
--
-- Two builds of one module can be in flight together — two merge plans a few minutes apart, each
-- asking for everything standing on what it changed — and they finish in any order. Each build
-- stands on the bases the mesh held when it was *asked*, so the one asked later is the newer one.
-- The mesh ordered builds by `at`, which is when the outcome was recorded, and registered whatever
-- it heard last: an older request that took longer replaced a newer one as what the module is, and
-- the next push sent machines an image built on a base the mesh had already replaced.
--
-- `build.asked` is when the build was requested, read from the correlation id the controller wrote
-- (`build-<unix nanoseconds>`). Nullable: an id of any other shape says no request time, and such a
-- build is placed where it was recorded, which is the order the mesh had before this.
alter table build add column asked timestamptz;
update build
set asked = to_timestamp((substring(id from '^build-([0-9]{19})$'))::numeric / 1000000000)
where id ~ '^build-[0-9]{19}$';
-- `module.built_asked` is when the build the module's registered manifest came from was asked, so
-- a later-heard outcome of an earlier request is recorded and not registered. A manifest handed over
-- by hand is a request made when it is handed over. Null for a module registered before this was
-- kept: its next registration, whichever it is, sets it.
alter table module add column built_asked timestamptz;
-124
View File
@@ -1,124 +0,0 @@
package inventory
import (
"context"
"errors"
"testing"
"time"
"github.com/novox/mesh-controller/internal/catalogue"
)
// novox/hq 04-ISSUES/219: two builds of one module in flight together, the one asked first heard
// last. The newer request stood on the newer base; the older one's late outcome is recorded and is
// not what the module is.
func postgresBuild(id string, asked time.Time, image string) Build {
b := aBuild(id, "postgres", "")
b.Asked = asked
b.Against = []string{"mesh-tools/runtime@sha256:" + id}
b.Made = []Artifact{{Name: "server", Kind: "image", Reference: "postgres@sha256:" + image}}
return b
}
func TestAnOlderRequestFinishingLaterIsNotWhatTheModuleHolds(t *testing.T) {
inv := fresh(t)
ctx := context.Background()
older := time.Date(2026, 10, 3, 21, 33, 45, 0, time.UTC)
newer := time.Date(2026, 10, 3, 21, 51, 57, 0, time.UTC)
// The newer request finishes first, the older one last — recorded in that order.
if err := inv.RecordBuild(ctx, postgresBuild("newer", newer, "4bcd5f73")); err != nil {
t.Fatal(err)
}
if err := inv.RecordBuild(ctx, postgresBuild("older", older, "0ab07fa9")); err != nil {
t.Fatal(err)
}
held, err := inv.Held(ctx)
if err != nil {
t.Fatal(err)
}
if got := held["postgres/server"]; got != "postgres@sha256:4bcd5f73" {
t.Errorf("postgres holds %q; want the newer request's image 4bcd5f73", got)
}
against, err := inv.BuiltAgainst(ctx)
if err != nil {
t.Fatal(err)
}
if got := against["postgres"]; len(got) != 1 || got[0] != "mesh-tools/runtime@sha256:newer" {
t.Errorf("postgres stands on %v; want what the newer request stood on", got)
}
// Both are still recorded, the late one first as what happened lately.
builds, err := inv.Builds(ctx, "postgres", 5)
if err != nil {
t.Fatal(err)
}
if len(builds) != 2 || builds[0].ID != "older" || !builds[0].Asked.Equal(older) {
t.Fatalf("both builds, newest heard first, with when they were asked: %+v", builds)
}
}
func TestABuildWithNoKnownRequestTimeIsOrderedByWhenItWasRecorded(t *testing.T) {
// What the mesh did before it kept the request time, so a row from before still answers.
inv := fresh(t)
ctx := context.Background()
for _, id := range []string{"first", "second"} {
b := aBuild(id, "shell", "")
b.Made = []Artifact{{Name: "config", Kind: "archive", Reference: "…/" + id}}
if err := inv.RecordBuild(ctx, b); err != nil {
t.Fatal(err)
}
}
held, err := inv.Held(ctx)
if err != nil {
t.Fatal(err)
}
if got := held["shell/config"]; got != "…/second" {
t.Errorf("shell holds %q; want the one recorded last", got)
}
}
func TestARegistrationFromAnOlderRequestDoesNotReplaceANewerOne(t *testing.T) {
inv := fresh(t)
ctx := context.Background()
older := time.Date(2026, 10, 3, 21, 33, 45, 0, time.UTC)
newer := time.Date(2026, 10, 3, 21, 51, 57, 0, time.UTC)
from := func(asked time.Time) Source {
return Source{Repository: "novox/mesh-catalog", Seat: "git", Path: "modules/postgres",
BuiltFrom: "efff5415", Asked: asked}
}
if err := inv.RegisterModule(ctx, catalogue.Manifest{Module: "postgres", Version: "fixed"}, from(newer)); err != nil {
t.Fatal(err)
}
err := inv.RegisterModule(ctx, catalogue.Manifest{Module: "postgres", Version: "stale"}, from(older))
if !errors.Is(err, ErrSuperseded) {
t.Fatalf("an older request's registration was not refused as superseded: %v", err)
}
shelf, err := inv.Catalogue(ctx)
if err != nil {
t.Fatal(err)
}
if got := shelf["postgres"].Version; got != "fixed" {
t.Fatalf("postgres is %q; want the newer request's manifest", got)
}
// A later request, and a manifest handed over by hand — asked when it is handed over — both
// replace it as before.
if err := inv.RegisterModule(ctx, catalogue.Manifest{Module: "postgres", Version: "later"},
from(newer.Add(time.Minute))); err != nil {
t.Fatal(err)
}
if err := inv.RegisterModule(ctx, catalogue.Manifest{Module: "postgres", Version: "by-hand"}, Source{}); err != nil {
t.Fatal(err)
}
if shelf, _ := inv.Catalogue(ctx); shelf["postgres"].Version != "by-hand" {
t.Fatalf("postgres is %q; want the manifest handed over by hand", shelf["postgres"].Version)
}
src, err := inv.SourceOf(ctx, "postgres")
if err != nil || src.Repository != "novox/mesh-catalog" {
t.Fatalf("a hand registration erased the provenance: %+v %v", src, err)
}
}
+2 -9
View File
@@ -44,15 +44,8 @@ func TestAPersonMayCallToolsAndNothingElse(t *testing.T) {
} }
// The one tool, both ways it is addressed (novox/hq ADR 0159): to whichever instance // The one tool, both ways it is addressed (novox/hq ADR 0159): to whichever instance
// answers, and to the instance on one machine. Nothing else. // answers, and to the instance on one machine. Nothing else.
// And asking what answers (novox/hq ADR 0197), which claims nothing and calls nothing. if len(perms.Publish) != 2 || perms.Publish[0] != "mesh.mod.mesh-catalog.tool.catalog_tools" ||
var tools []string perms.Publish[1] != "mesh.mod.mesh-catalog.tool.catalog_tools.*" {
for _, s := range perms.Publish {
if !strings.HasPrefix(s, "$SRV.") {
tools = append(tools, s)
}
}
if len(tools) != 2 || tools[0] != "mesh.mod.mesh-catalog.tool.catalog_tools" ||
tools[1] != "mesh.mod.mesh-catalog.tool.catalog_tools.*" {
t.Errorf("ada may publish %v, which should be the one tool, both ways addressed, and nothing else", perms.Publish) t.Errorf("ada may publish %v, which should be the one tool, both ways addressed, and nothing else", perms.Publish)
} }
for _, s := range perms.Publish { for _, s := range perms.Publish {
+2 -15
View File
@@ -17,7 +17,7 @@ import (
// DiscoverySubjects are where one service instance is asked to say what it is. // DiscoverySubjects are where one service instance is asked to say what it is.
func DiscoverySubjects(name, id string) []string { func DiscoverySubjects(name, id string) []string {
var out []string var out []string
for _, verb := range []string{"PING", "INFO", "STATS"} { for _, verb := range []string{"PING", "INFO"} {
out = append(out, "$SRV."+verb, "$SRV."+verb+"."+name, "$SRV."+verb+"."+name+"."+id) out = append(out, "$SRV."+verb, "$SRV."+verb+"."+name, "$SRV."+verb+"."+name+"."+id)
} }
return out return out
@@ -35,16 +35,6 @@ func (b OverNATS) Announce(info micro.Info, logger *log.Logger) (func(), error)
if err != nil { if err != nil {
return nil, err return nil, err
} }
// Statistics the protocol asks for; the controller keeps none per verb, so it answers its
// identity and its endpoints with nothing counted — an honest zero, not a refusal.
stats := micro.Stats{ServiceIdentity: info.ServiceIdentity, Type: micro.StatsResponseType}
for _, e := range info.Endpoints {
stats.Endpoints = append(stats.Endpoints, &micro.EndpointStats{Name: e.Name, Subject: e.Subject, QueueGroup: e.QueueGroup})
}
statsBody, err := json.Marshal(stats)
if err != nil {
return nil, err
}
var subs []*nats.Subscription var subs []*nats.Subscription
done := make(chan struct{}) done := make(chan struct{})
stop := func() { stop := func() {
@@ -56,11 +46,8 @@ func (b OverNATS) Announce(info micro.Info, logger *log.Logger) (func(), error)
for _, subject := range DiscoverySubjects(info.Name, info.ID) { for _, subject := range DiscoverySubjects(info.Name, info.ID) {
subject := subject subject := subject
body := infoBody body := infoBody
switch { if len(subject) >= 9 && subject[:9] == "$SRV.PING" {
case len(subject) >= 9 && subject[:9] == "$SRV.PING":
body = pingBody body = pingBody
case len(subject) >= 10 && subject[:10] == "$SRV.STATS":
body = statsBody
} }
bind := func() (*nats.Subscription, error) { bind := func() (*nats.Subscription, error) {
return b.Conn.Subscribe(subject, func(msg *nats.Msg) { return b.Conn.Subscribe(subject, func(msg *nats.Msg) {
-29
View File
@@ -2,9 +2,6 @@ package link
import ( import (
"encoding/json" "encoding/json"
"strconv"
"strings"
"time"
) )
// Asking a machine to build a module, and hearing what came out. // Asking a machine to build a module, and hearing what came out.
@@ -20,32 +17,6 @@ import (
// holds no opinion about what they contain, and a host that also built things would be a host // holds no opinion about what they contain, and a host that also built things would be a host
// with a container runtime requirement and a git dependency (novox/hq ADR 0005). // with a container runtime requirement and a git dependency (novox/hq ADR 0005).
// NewBuildID is the correlation for a build asked at that moment: `build-<unix nanoseconds>`.
//
// **The id carries when the build was asked, and that is read back** (novox/hq 04-ISSUES/219). Builds
// of one module can be in flight together and finish in any order; what a module currently is must
// be the newest *request's* outcome, not the last one heard, and the id is the one thing every
// outcome echoes whichever builder answered it. One place writes the shape and one reads it.
func NewBuildID(asked time.Time) string {
return "build-" + strconv.FormatInt(asked.UnixNano(), 10)
}
// BuildAskedAt is when the build with this id was asked, as NewBuildID wrote it. False for an id
// of any other shape — one written before this was read, or by hand — whose request time the mesh
// does not know.
func BuildAskedAt(id string) (time.Time, bool) {
digits, ok := strings.CutPrefix(id, "build-")
if !ok || digits == "" {
return time.Time{}, false
}
nanos, err := strconv.ParseInt(digits, 10, 64)
// A number too small to be a moment this mesh could have asked at is a name, not a time.
if err != nil || nanos < time.Date(2020, 1, 1, 0, 0, 0, 0, time.UTC).UnixNano() {
return time.Time{}, false
}
return time.Unix(0, nanos).UTC(), true
}
// BuildRequest is one module to build. // BuildRequest is one module to build.
type BuildRequest struct { type BuildRequest struct {
// ID correlates the answer with the asking. Not the module name: two builds of one module can // ID correlates the answer with the asking. Not the module name: two builds of one module can
+2 -54
View File
@@ -5,7 +5,6 @@ import (
"encoding/json" "encoding/json"
"errors" "errors"
"fmt" "fmt"
"log"
"strings" "strings"
"time" "time"
@@ -79,15 +78,10 @@ func (n *natsInbound) Receive(ctx context.Context, act func(context.Context, Con
// than creating one here: the consumer is an object with a configuration — ack policy, ack // than creating one here: the consumer is an object with a configuration — ack policy, ack
// wait, redelivery — and a client that creates its own would be a second opinion about it. // wait, redelivery — and a client that creates its own would be a second opinion about it.
control := make(chan *nats.Msg, Prefetch) control := make(chan *nats.Msg, Prefetch)
said, err := standingBy(ctx, log.Default(), "CONTROL", func() (*nats.Subscription, error) { said, err := js.ChanSubscribe("", control, nats.Bind("CONTROL", broker.ControllerName))
return js.ChanSubscribe("", control, nats.Bind("CONTROL", broker.ControllerName))
})
if err != nil { if err != nil {
return fmt.Errorf("subscribing to what nodes say: %w", err) return fmt.Errorf("subscribing to what nodes say: %w", err)
} }
if said == nil {
return nil // stopped while standing by
}
defer func() { _ = said.Unsubscribe() }() defer func() { _ = said.Unsubscribe() }()
// Heartbeats, on core NATS and off any stream (design 25 §3). Their own subscription because // Heartbeats, on core NATS and off any stream (design 25 §3). Their own subscription because
@@ -103,15 +97,10 @@ func (n *natsInbound) Receive(ctx context.Context, act func(context.Context, Con
var events chan *nats.Msg var events chan *nats.Msg
if len(n.follows) > 0 { if len(n.follows) > 0 {
events = make(chan *nats.Msg, Prefetch) events = make(chan *nats.Msg, Prefetch)
followed, err := standingBy(ctx, log.Default(), "EVENTS", func() (*nats.Subscription, error) { followed, err := js.ChanSubscribe("", events, nats.Bind("EVENTS", broker.ControllerName))
return js.ChanSubscribe("", events, nats.Bind("EVENTS", broker.ControllerName))
})
if err != nil { if err != nil {
return fmt.Errorf("subscribing to what the catalogue says: %w", err) return fmt.Errorf("subscribing to what the catalogue says: %w", err)
} }
if followed == nil {
return nil
}
defer func() { _ = followed.Unsubscribe() }() defer func() { _ = followed.Unsubscribe() }()
} }
@@ -335,44 +324,3 @@ func (m *natsControl) forget() {
type replyAddressed struct { type replyAddressed struct {
ReplyTo string `json:"reply_to,omitempty"` ReplyTo string `json:"reply_to,omitempty"`
} }
// StandbyPoll is how often a controller standing by looks again for its consumers. A variable so a
// test need not wait.
var StandbyPoll = 2 * time.Second
// standingBy binds one of the controller's consumers, waiting while another controller holds it.
//
// **Two controllers, one consumer** (novox/hq issue 213). The controller's consumers are push
// consumers with no delivery group, so the server lets one subscription bind each — on purpose:
// two would each act on every message (issue 146). When a machine hands its controller over from
// the container to the process, the host starts the process first and removes the container only
// once the process is up; the process then finds the consumers bound. Exiting on that would never
// be up, so the container would never go. It stands by instead — the seat's verbs are already
// served from a queue group, and the plans wait on their lock — and binds as soon as the other lets
// go. Nil and no error is ctx ending while it waited.
func standingBy(ctx context.Context, logger interface{ Printf(string, ...any) }, stream string,
bind func() (*nats.Subscription, error)) (*nats.Subscription, error) {
said := false
for {
sub, err := bind()
if err == nil {
if said {
logger.Printf("took the controller's consumer on %s: the controller that held it let go", stream)
}
return sub, nil
}
if !strings.Contains(err.Error(), "already bound") {
return nil, err
}
if !said {
logger.Printf("another controller holds the controller's consumer on %s; standing by "+
"until it lets go", stream)
said = true
}
select {
case <-ctx.Done():
return nil, nil
case <-time.After(StandbyPoll):
}
}
}
-69
View File
@@ -1,69 +0,0 @@
package link
import (
"context"
"encoding/json"
"os"
"testing"
"time"
"github.com/novox/mesh-controller/internal/broker"
)
// novox/hq issue 213: while a machine hands its controller over, the new controller (the process)
// is started while the old one (the container) still holds the controller's consumers. It must not
// exit — the host would read that as a replacement that did not come up and never remove the
// container — and must not act on what the old one is handed. It stands by, and takes the consumers
// when the old one lets go.
func TestNatsASecondControllerStandsByAndTakesOverWhenTheFirstLetsGo(t *testing.T) {
js := aBus(t)
was := StandbyPoll
StandbyPoll = 50 * time.Millisecond
defer func() { StandbyPoll = was }()
old := &counted{}
_, stopOld := servingOn(t, js, old)
eventually(t, "the first controller binding its consumer", func() bool {
info, err := js.Context().ConsumerInfo("CONTROL", broker.ControllerName)
return err == nil && info.PushBound
})
// The new one, on a connection of its own as the process would have.
second, err := broker.Dial(os.Getenv("MESH_TEST_NATS"))
if err != nil {
t.Fatal(err)
}
t.Cleanup(second.Close)
fresh := &counted{}
s := &Server{inbound: Nats(second), bus: OverNATS{Conn: second.Conn(), JS: second.Context()},
listener: fresh, log: quiet()}
ctx, stopNew := context.WithCancel(context.Background())
defer stopNew()
ended := make(chan error, 1)
go func() { ended <- s.Serve(ctx) }()
select {
case err := <-ended:
t.Fatalf("the second controller stopped instead of standing by: %v", err)
case <-time.After(500 * time.Millisecond):
}
report := func(declared string) {
body, _ := json.Marshal(Report{Node: "anchor", Declared: declared, Applied: []string{"store"}})
if _, err := js.Context().Publish(ReportSubject("anchor"), body); err != nil {
t.Fatal(err)
}
}
report("d1")
eventually(t, "the holding controller hearing the report", func() bool { return old.count() == 1 })
if fresh.count() != 0 {
t.Fatal("the controller standing by acted on a report the holder was handed")
}
stopOld()
report("d2")
eventually(t, "the second controller taking over once the first let go", func() bool { return fresh.count() == 1 })
if old.count() != 1 {
t.Errorf("the first controller heard %d reports", old.count())
}
}