Compare commits
17
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
da394b45e6 | ||
|
|
2f3bfda8c0 | ||
|
|
208388978a | ||
|
|
aa771616bb | ||
|
|
1513bbaac9 | ||
|
|
0014984116 | ||
|
|
d6e49dbd68 | ||
|
|
3756bb3460 | ||
|
|
60be9c5360 | ||
|
|
da31bcb11e | ||
|
|
f03e7b33c9 | ||
|
|
0baf727f36 | ||
|
|
94dd49a968 | ||
|
|
83a298e7e0 | ||
|
|
220b79f5cd | ||
|
|
e62201e227 | ||
|
|
0ab9b86f0a |
@@ -194,6 +194,9 @@ func answer(ctx context.Context, publisher builder.Publisher, on, workspace stri
|
|||||||
})
|
})
|
||||||
}
|
}
|
||||||
result.Against = built.Against
|
result.Against = built.Against
|
||||||
|
for _, r := range built.Read {
|
||||||
|
result.Read = append(result.Read, link.ReadRepository{Repository: r.Repository, Ref: r.Ref})
|
||||||
|
}
|
||||||
fmt.Fprintf(os.Stderr, " built %s from %s\n", built.Manifest.Module, short(built.Commit))
|
fmt.Fprintf(os.Stderr, " built %s from %s\n", built.Manifest.Module, short(built.Commit))
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -106,6 +106,9 @@ func buildOnce(ctx context.Context, args []string) error {
|
|||||||
Manifest: built.Manifest,
|
Manifest: built.Manifest,
|
||||||
Against: built.Against,
|
Against: built.Against,
|
||||||
}
|
}
|
||||||
|
for _, r := range built.Read {
|
||||||
|
out.Read = append(out.Read, readRepository{Repository: r.Repository, Ref: r.Ref})
|
||||||
|
}
|
||||||
for _, made := range built.Built {
|
for _, made := range built.Built {
|
||||||
out.Made = append(out.Made, madeArtifact{Name: made.Name, Kind: made.Kind, Reference: made.Reference})
|
out.Made = append(out.Made, madeArtifact{Name: made.Name, Kind: made.Kind, Reference: made.Reference})
|
||||||
}
|
}
|
||||||
@@ -131,6 +134,13 @@ type onceResult struct {
|
|||||||
Manifest any `json:"manifest"`
|
Manifest any `json:"manifest"`
|
||||||
Made []madeArtifact `json:"made"`
|
Made []madeArtifact `json:"made"`
|
||||||
Against []string `json:"against,omitempty"`
|
Against []string `json:"against,omitempty"`
|
||||||
|
Read []readRepository `json:"read,omitempty"`
|
||||||
|
}
|
||||||
|
|
||||||
|
// readRepository is a repository this build read source from besides the module's own.
|
||||||
|
type readRepository struct {
|
||||||
|
Repository string `json:"repository"`
|
||||||
|
Ref string `json:"ref,omitempty"`
|
||||||
}
|
}
|
||||||
|
|
||||||
type madeArtifact struct {
|
type madeArtifact struct {
|
||||||
|
|||||||
@@ -149,6 +149,9 @@ func buildFrom(result link.BuildResult) inventory.Build {
|
|||||||
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))
|
||||||
}
|
}
|
||||||
|
for _, r := range result.Read {
|
||||||
|
kept.Read = append(kept.Read, inventory.ReadRepository{Repository: r.Repository, Ref: r.Ref})
|
||||||
|
}
|
||||||
var announced []inventory.Artifact
|
var announced []inventory.Artifact
|
||||||
for _, made := range result.Made {
|
for _, made := range result.Made {
|
||||||
announced = append(announced, inventory.Artifact{
|
announced = append(announced, inventory.Artifact{
|
||||||
|
|||||||
@@ -69,12 +69,20 @@ func moduleCommand(ctx context.Context, args []string) error {
|
|||||||
repo := set.String("source", "", "where this module comes from")
|
repo := set.String("source", "", "where this module comes from")
|
||||||
ref := set.String("ref", "", "the branch followed there")
|
ref := set.String("ref", "", "the branch followed there")
|
||||||
commit := set.String("commit", "", "the commit this manifest was read at")
|
commit := set.String("commit", "", "the commit this manifest was read at")
|
||||||
|
// **Where inside the repository the module is** (novox/hq ADR 0069). A module is a
|
||||||
|
// repository *and* a directory, and a record that carries only the repository names a
|
||||||
|
// module.json at its root — so every later build of it looks in the wrong place and fails
|
||||||
|
// with "no module.json at its root". Nine modules on this mesh were registered that way
|
||||||
|
// and none of them could be rebuilt (2026-09-28).
|
||||||
|
path := set.String("path", "", "the module's directory inside that repository")
|
||||||
|
self := set.Bool("self", false, "the source is a path on the forge holding the git seat")
|
||||||
positionals, err := parseAround(set, args[1:])
|
positionals, err := parseAround(set, args[1:])
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
if len(positionals) != 1 {
|
if len(positionals) != 1 {
|
||||||
return errors.New("module add <manifest.json> [--source <repo> --ref <branch> --commit <sha>]")
|
return errors.New("module add <manifest.json> [--source <repo> [--self] [--path P] " +
|
||||||
|
"--ref <branch> --commit <sha>]")
|
||||||
}
|
}
|
||||||
raw, err := os.ReadFile(positionals[0])
|
raw, err := os.ReadFile(positionals[0])
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@@ -84,23 +92,24 @@ func moduleCommand(ctx context.Context, args []string) error {
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
// Provenance together or not at all. A source with no commit cannot be compared against
|
from, err := whereItComesFrom(*repo, *ref, *commit, *path, *self)
|
||||||
// anything, so it would record where the module came from and still never be able to say
|
if err != nil {
|
||||||
// the mesh is behind it — which is the one thing recording it is for.
|
return err
|
||||||
if (*repo == "") != (*commit == "") {
|
|
||||||
return errors.New("--source and --commit go together: a source with no commit " +
|
|
||||||
"cannot be compared against anything, and a commit with no source has nothing " +
|
|
||||||
"to be compared with")
|
|
||||||
}
|
}
|
||||||
if err := inv.RegisterModule(ctx, m, inventory.Source{
|
if err := inv.RegisterModule(ctx, m, from); err != nil {
|
||||||
Repository: *repo, Ref: *ref, BuiltFrom: *commit,
|
|
||||||
}); err != nil {
|
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
fmt.Printf("%s registered", m.Module)
|
fmt.Printf("%s registered", m.Module)
|
||||||
if *commit != "" {
|
if *commit != "" {
|
||||||
fmt.Printf(" from %s", short(*commit))
|
fmt.Printf(" from %s", short(*commit))
|
||||||
}
|
}
|
||||||
|
if *repo != "" && *path == "" {
|
||||||
|
// Said, not refused: a module really at the root is the ordinary case for a repository
|
||||||
|
// of its own. But a repository holding many modules and a record naming none of them is
|
||||||
|
// a module nothing can rebuild, and the person adding it is the one who knows which.
|
||||||
|
fmt.Printf("\n no directory inside %s, so it is built from that repository's root — "+
|
||||||
|
"`--path` if the module lives in a directory there", *repo)
|
||||||
|
}
|
||||||
if len(m.Provides) > 0 {
|
if len(m.Provides) > 0 {
|
||||||
fmt.Printf(", providing %s", describeOffers(m.Provides))
|
fmt.Printf(", providing %s", describeOffers(m.Provides))
|
||||||
}
|
}
|
||||||
@@ -602,3 +611,37 @@ func issueWith(ctx context.Context, inv *inventory.Inventory, m catalogue.Manife
|
|||||||
"machine holding mesh-broker\n")
|
"machine holding mesh-broker\n")
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// whereItComesFrom is the provenance a module handed over by hand records, and what a record must
|
||||||
|
// say to be worth anything later.
|
||||||
|
//
|
||||||
|
// **A module is a repository and a directory inside it** (novox/hq ADR 0069). A record carrying only
|
||||||
|
// the repository names a module.json at its root, so every later build of it looks in the wrong
|
||||||
|
// place — nine modules on this mesh were registered that way and none of them could be rebuilt
|
||||||
|
// (2026-09-28). The directory cannot be checked from here, because the control plane does not clone;
|
||||||
|
// what can be checked is that the record is whole.
|
||||||
|
func whereItComesFrom(repository, ref, commit, path string, self bool) (inventory.Source, error) {
|
||||||
|
// Provenance together or not at all. A source with no commit cannot be compared against
|
||||||
|
// anything, so it would record where the module came from and still never be able to say the
|
||||||
|
// mesh is behind it — which is the one thing recording it is for.
|
||||||
|
if (repository == "") != (commit == "") {
|
||||||
|
return inventory.Source{}, errors.New("--source and --commit go together: a source with " +
|
||||||
|
"no commit cannot be compared against anything, and a commit with no source has " +
|
||||||
|
"nothing to be compared with")
|
||||||
|
}
|
||||||
|
// A directory or a forge with no repository is half a location, and the half it keeps is the
|
||||||
|
// half nothing can be found with.
|
||||||
|
if repository == "" && (path != "" || self) {
|
||||||
|
return inventory.Source{}, errors.New("--path and --self say where inside a source and " +
|
||||||
|
"which forge holds it, so they need --source: without one there is nothing for them " +
|
||||||
|
"to be part of")
|
||||||
|
}
|
||||||
|
from := inventory.Source{Repository: repository, Ref: ref, BuiltFrom: commit, Path: path}
|
||||||
|
if self {
|
||||||
|
if err := onASeat(repository); err != nil {
|
||||||
|
return inventory.Source{}, err
|
||||||
|
}
|
||||||
|
from.Seat = gitSeat
|
||||||
|
}
|
||||||
|
return from, nil
|
||||||
|
}
|
||||||
|
|||||||
@@ -189,13 +189,17 @@ func readPrivateKey(path string) (string, error) {
|
|||||||
func personIssue(ctx context.Context, args []string) error {
|
func personIssue(ctx context.Context, args []string) error {
|
||||||
set := flag.NewFlagSet("operator issue", flag.ContinueOnError)
|
set := flag.NewFlagSet("operator issue", flag.ContinueOnError)
|
||||||
invokes := set.String("invokes", "", "the tools this person may call, comma-separated, or * for every one")
|
invokes := set.String("invokes", "", "the tools this person may call, comma-separated, or * for every one")
|
||||||
if err := set.Parse(args); err != nil {
|
// Flags on either side of the name, because the usage this command prints puts them after it —
|
||||||
|
// and the standard parser stops at the first thing that is not a flag, so the order the command
|
||||||
|
// documents was the one order it refused (2026-09-28).
|
||||||
|
positionals, err := parseAround(set, args)
|
||||||
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
if set.NArg() != 1 {
|
if len(positionals) != 1 {
|
||||||
return errors.New("operator issue <name> --invokes <tool,tool|*>")
|
return errors.New("operator issue <name> --invokes <tool,tool|*>")
|
||||||
}
|
}
|
||||||
name := set.Arg(0)
|
name := positionals[0]
|
||||||
if *invokes == "" {
|
if *invokes == "" {
|
||||||
return errors.New(
|
return errors.New(
|
||||||
"say what this person may call: --invokes mesh-catalog.catalog_tools,gitea.repo_create, " +
|
"say what this person may call: --invokes mesh-catalog.catalog_tools,gitea.repo_create, " +
|
||||||
|
|||||||
@@ -3,6 +3,7 @@ package main
|
|||||||
import (
|
import (
|
||||||
"strings"
|
"strings"
|
||||||
"testing"
|
"testing"
|
||||||
|
"time"
|
||||||
|
|
||||||
"github.com/novox/mesh-controller/internal/catalogue"
|
"github.com/novox/mesh-controller/internal/catalogue"
|
||||||
"github.com/novox/mesh-controller/internal/inventory"
|
"github.com/novox/mesh-controller/internal/inventory"
|
||||||
@@ -90,3 +91,122 @@ func TestAMergeMatchesTheSourcesBuiltFromIt(t *testing.T) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// A merge made before the source was last seen is history: it does not move the source, and a
|
||||||
|
// merge that says nothing about when it was made is taken as news.
|
||||||
|
func TestAMergeOlderThanTheLastLookIsHistory(t *testing.T) {
|
||||||
|
seen := time.Date(2026, 9, 28, 3, 0, 0, 0, time.UTC)
|
||||||
|
if !isHistory("2026-09-28T02:00:00Z", seen) {
|
||||||
|
t.Fatal("an older merge was taken as news")
|
||||||
|
}
|
||||||
|
if isHistory("2026-09-28T04:00:00Z", seen) {
|
||||||
|
t.Fatal("a newer merge was taken as history")
|
||||||
|
}
|
||||||
|
if isHistory("", seen) || isHistory("2026-09-28T02:00:00Z", time.Time{}) {
|
||||||
|
t.Fatal("a merge or a source with no time on it was refused")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// A module as the catalogue holds it: built from a repository, at a directory inside it.
|
||||||
|
func fromRepo(module, repository, path string) inventory.Entry {
|
||||||
|
return inventory.Entry{
|
||||||
|
Manifest: catalogue.Manifest{Module: module},
|
||||||
|
Source: inventory.Source{Repository: repository, Path: path, Ref: "main"},
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// A merge rebuilds the modules whose own directories it changed, and everything when what it changed
|
||||||
|
// is shared. One repository holding many modules is the ordinary case here, and rebuilding all of
|
||||||
|
// them for a change to one is what exhausted a registry's pull limit the first night this ran.
|
||||||
|
func TestAMergeRebuildsTheModulesItChanged(t *testing.T) {
|
||||||
|
const repo = "http://forge.internal:20000/novox/mesh-catalog.git"
|
||||||
|
gitea := fromRepo("gitea", repo, "modules/gitea")
|
||||||
|
keycloak := fromRepo("keycloak", repo, "modules/keycloak")
|
||||||
|
known := []inventory.Entry{gitea, keycloak, fromRepo("plex", repo, "modules/plex")}
|
||||||
|
candidates := []inventory.Entry{gitea, keycloak}
|
||||||
|
merge := func(paths []string, truncated bool) link.SourceMoved {
|
||||||
|
return link.SourceMoved{Owner: "novox", Repo: "mesh-catalog", Base: "main",
|
||||||
|
Paths: paths, PathsTruncated: truncated}
|
||||||
|
}
|
||||||
|
named := func(entries []inventory.Entry) string {
|
||||||
|
var names []string
|
||||||
|
for _, e := range entries {
|
||||||
|
names = append(names, e.Manifest.Module)
|
||||||
|
}
|
||||||
|
return strings.Join(names, ",")
|
||||||
|
}
|
||||||
|
for _, c := range []struct {
|
||||||
|
what string
|
||||||
|
m link.SourceMoved
|
||||||
|
want string
|
||||||
|
}{
|
||||||
|
{"one module's own files", merge([]string{"modules/gitea/index.ts", "modules/gitea/client.ts"}, false), "gitea"},
|
||||||
|
{"two modules' files", merge([]string{"modules/gitea/index.ts", "modules/keycloak/module.json"}, false), "gitea,keycloak"},
|
||||||
|
{"a file they share", merge([]string{"tsconfig.json"}, false), "gitea,keycloak"},
|
||||||
|
{"a module the mesh does not hold", merge([]string{"modules/plex/index.ts"}, false), ""},
|
||||||
|
{"nothing said about the files", merge(nil, false), "gitea,keycloak"},
|
||||||
|
{"more files than were listed", merge([]string{"modules/gitea/index.ts"}, true), "gitea,keycloak"},
|
||||||
|
} {
|
||||||
|
if got := named(whatTheMergeTouched(candidates, known, c.m)); got != c.want {
|
||||||
|
t.Errorf("%s: rebuilt %q, wanted %q", c.what, got, c.want)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// A module whose recipe packages source from another repository is affected when that repository
|
||||||
|
// moves — the manifest the mesh keeps says nothing about it, so the record of what the build read is
|
||||||
|
// the only thing that can say so.
|
||||||
|
func TestAModuleIsAffectedByTheRepositoryItPackages(t *testing.T) {
|
||||||
|
m := link.SourceMoved{Owner: "novox", Repo: "mesh-controller", Base: "main",
|
||||||
|
CloneURL: "http://forge.internal:20000/novox/mesh-controller.git"}
|
||||||
|
for _, read := range [][]inventory.ReadRepository{
|
||||||
|
{{Repository: "http://forge.internal:20000/novox/mesh-controller.git", Ref: "main"}},
|
||||||
|
{{Repository: "novox/mesh-controller"}},
|
||||||
|
{{Repository: "https://elsewhere.example/novox/other"}, {Repository: "novox/mesh-controller.git", Ref: "main"}},
|
||||||
|
} {
|
||||||
|
if !readsFrom(read, m) {
|
||||||
|
t.Errorf("%+v was not matched by the merge", read)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
for _, read := range [][]inventory.ReadRepository{
|
||||||
|
nil,
|
||||||
|
{{Repository: "novox/mesh-host", Ref: "main"}},
|
||||||
|
{{Repository: "novox/mesh-controller", Ref: "release"}},
|
||||||
|
} {
|
||||||
|
if readsFrom(read, m) {
|
||||||
|
t.Errorf("%+v was matched by a merge that is not its", read)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// What a module handed over by hand records about where it came from, and what is refused.
|
||||||
|
func TestWhatAHandedOverModuleRecordsAboutItsSource(t *testing.T) {
|
||||||
|
// The whole location: a repository on the mesh's own forge, the directory inside it, the branch
|
||||||
|
// and the commit the manifest was read at.
|
||||||
|
from, err := whereItComesFrom("novox/mesh-catalog", "main", "c0ffee", "modules/gitea", true)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if from.Path != "modules/gitea" || from.Seat != "git" || from.Repository != "novox/mesh-catalog" {
|
||||||
|
t.Fatalf("the source records as %+v", from)
|
||||||
|
}
|
||||||
|
// A manifest with no provenance at all is legitimate: fixing something in a hurry.
|
||||||
|
if from, err := whereItComesFrom("", "", "", "", false); err != nil || from != (inventory.Source{}) {
|
||||||
|
t.Fatalf("a manifest handed over with no provenance was refused: %+v, %v", from, err)
|
||||||
|
}
|
||||||
|
for _, c := range []struct {
|
||||||
|
what string
|
||||||
|
repository, ref, commit, path string
|
||||||
|
self bool
|
||||||
|
}{
|
||||||
|
{what: "a source with no commit", repository: "novox/mesh-catalog", commit: ""},
|
||||||
|
{what: "a commit with no source", commit: "c0ffee"},
|
||||||
|
{what: "a directory inside nothing", path: "modules/gitea"},
|
||||||
|
{what: "a forge holding nothing", self: true},
|
||||||
|
{what: "an address given as a path on the forge", repository: "http://forge.internal:20000/novox/x.git", commit: "c0ffee", self: true},
|
||||||
|
} {
|
||||||
|
if _, err := whereItComesFrom(c.repository, c.ref, c.commit, c.path, c.self); err == nil {
|
||||||
|
t.Errorf("%s was recorded as a source", c.what)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -753,8 +753,31 @@ func raiseTheBus(ctx context.Context, inv *inventory.Inventory, address string)
|
|||||||
if err := broker.RaiseSeats(js, inventory.MeshSeats(), holders); err != nil {
|
if err := broker.RaiseSeats(js, inventory.MeshSeats(), holders); err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
fmt.Printf("the bus at %s has its streams, and %d machine(s) can hear a declaration\n",
|
// And how every module hears what it consumes. Derived from the same records the user list is
|
||||||
broker.BareAddress(address), len(names))
|
// composed from, so a module the mesh grants a consumer's subjects has that consumer waiting.
|
||||||
|
// Done on every raise, not only when a credential is issued: every module moved onto this bus
|
||||||
|
// by the rollout was issued on the old one, and came up with nothing to bind to (2026-09-28).
|
||||||
|
records, err := inv.BusRecords(ctx)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
users, err := broker.Users(records)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
hearing := 0
|
||||||
|
for _, p := range users {
|
||||||
|
consumer, needed := broker.ConsumerFor(p)
|
||||||
|
if !needed {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
if err := js.EnsureConsumer(consumer); err != nil {
|
||||||
|
return fmt.Errorf("how %s on %s hears what it consumes: %w", p.Module, p.Node, err)
|
||||||
|
}
|
||||||
|
hearing++
|
||||||
|
}
|
||||||
|
fmt.Printf("the bus at %s has its streams, %d machine(s) can hear a declaration, and %d module(s) "+
|
||||||
|
"can hear what they consume\n", broker.BareAddress(address), len(names), hearing)
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -127,9 +127,13 @@ func readinessOf(ctx context.Context, inv *inventory.Inventory) (broker.Readines
|
|||||||
if address != "" {
|
if address != "" {
|
||||||
// One dial, briefly. "Is it answering" is the one fact records cannot hold, and a mesh about
|
// One dial, briefly. "Is it answering" is the one fact records cannot hold, and a mesh about
|
||||||
// to move onto a server that is not there should hear it here rather than afterwards.
|
// to move onto a server that is not there should hear it here rather than afterwards.
|
||||||
if conn, err := nats.Connect(broker.BareAddress(address), nats.Timeout(5*time.Second)); err == nil {
|
//
|
||||||
|
// **Dialled the way the mesh dials it** — credential and pin — because a bare connect to a
|
||||||
|
// bus that requires TLS and a user fails at the handshake, and the check then reported a
|
||||||
|
// standing server as absent (seen live, 2026-09-28).
|
||||||
|
if js, err := broker.Dial(address, nats.Timeout(5*time.Second)); err == nil {
|
||||||
state.ServerStanding = true
|
state.ServerStanding = true
|
||||||
conn.Close()
|
js.Close()
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
+172
-12
@@ -232,22 +232,67 @@ func (f following) SourceMoved(ctx context.Context, m link.SourceMoved) error {
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
return notNow(err)
|
return notNow(err)
|
||||||
}
|
}
|
||||||
var moved []inventory.Entry
|
read, err := inv.ReadRepositories(ctx)
|
||||||
|
if err != nil {
|
||||||
|
return notNow(err)
|
||||||
|
}
|
||||||
|
|
||||||
|
// Two kinds of module are affected by one merge, and they are affected differently.
|
||||||
|
//
|
||||||
|
// A module **built from** this repository and branch has moved: the mesh records the new commit
|
||||||
|
// as what its source now has, and only what the merge actually changed is rebuilt. A module that
|
||||||
|
// only **packages source from** it has not moved — its own source is somewhere else, at the
|
||||||
|
// commit it already records — so it is rebuilt and its record left alone. Writing this commit as
|
||||||
|
// its source would make it permanently behind a repository its manifest does not come from.
|
||||||
|
var from, packaging []inventory.Entry
|
||||||
|
already := 0
|
||||||
for _, e := range entries {
|
for _, e := range entries {
|
||||||
if !sourceIs(e.Source, m) {
|
switch {
|
||||||
continue
|
case sourceIs(e.Source, m):
|
||||||
}
|
|
||||||
if e.Source.BuiltFrom == m.Commit {
|
if e.Source.BuiltFrom == m.Commit {
|
||||||
|
already++
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
|
// **A merge older than the last look at the source is history, not a move.** The forge
|
||||||
|
// announces what it finds merged, and an old merge surfacing late would otherwise move
|
||||||
|
// the recorded head backwards and rebuild everything built from that repository, once
|
||||||
|
// per old merge (2026-09-28).
|
||||||
|
if isHistory(m.MergedAt, e.Source.Seen) {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
from = append(from, e)
|
||||||
|
case readsFrom(read[e.Manifest.Module], m):
|
||||||
|
packaging = append(packaging, e)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if len(from) == 0 && len(packaging) == 0 {
|
||||||
|
// "Already built from it" and "nothing reads it" are different facts, and reading the first
|
||||||
|
// as the second sends somebody looking for a broken trigger when the mesh is up to date.
|
||||||
|
if already > 0 {
|
||||||
|
fmt.Printf("%s/%s merged into %s (%.8s); %d module(s) the mesh holds are already built "+
|
||||||
|
"from it\n", m.Owner, m.Repo, m.Base, m.Commit, already)
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
fmt.Printf("%s/%s merged into %s (%.8s); nothing the mesh holds reads it\n",
|
||||||
|
m.Owner, m.Repo, m.Base, m.Commit)
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
// The same judgement for the packaging kind, against the newest look at that repository by
|
||||||
|
// anything built from it: they keep no record of it themselves, and a replayed old merge should
|
||||||
|
// not rebuild them either.
|
||||||
|
if isHistory(m.MergedAt, lastLookAt(entries, m)) {
|
||||||
|
packaging = nil
|
||||||
|
}
|
||||||
|
touched := whatTheMergeTouched(from, entries, m)
|
||||||
|
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 {
|
||||||
return notNow(err)
|
return notNow(err)
|
||||||
}
|
}
|
||||||
moved = append(moved, e)
|
|
||||||
}
|
}
|
||||||
|
moved := append(append([]inventory.Entry{}, touched...), packaging...)
|
||||||
if len(moved) == 0 {
|
if len(moved) == 0 {
|
||||||
fmt.Printf("%s/%s merged into %s (%.8s); nothing the mesh holds is built from it\n",
|
fmt.Printf("%s/%s merged into %s (%.8s); it changed nothing any module the mesh holds is "+
|
||||||
m.Owner, m.Repo, m.Base, m.Commit)
|
"built from\n", m.Owner, m.Repo, m.Base, m.Commit)
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
against, err := inv.BuiltAgainst(ctx)
|
against, err := inv.BuiltAgainst(ctx)
|
||||||
@@ -261,6 +306,14 @@ func (f following) SourceMoved(ctx context.Context, m link.SourceMoved) error {
|
|||||||
}
|
}
|
||||||
fmt.Printf("%s/%s merged into %s (%.8s); building %s\n",
|
fmt.Printf("%s/%s merged into %s (%.8s); building %s\n",
|
||||||
m.Owner, m.Repo, m.Base, m.Commit, strings.Join(names, ", "))
|
m.Owner, m.Repo, m.Base, m.Commit, strings.Join(names, ", "))
|
||||||
|
if len(packaging) > 0 {
|
||||||
|
var also []string
|
||||||
|
for _, e := range packaging {
|
||||||
|
also = append(also, e.Manifest.Module)
|
||||||
|
}
|
||||||
|
fmt.Printf(" %s package source from it, so they are rebuilt and their own source record "+
|
||||||
|
"is left where it is\n", strings.Join(also, ", "))
|
||||||
|
}
|
||||||
var failed []string
|
var failed []string
|
||||||
for _, e := range ordered {
|
for _, e := range ordered {
|
||||||
source := buildSource{Repository: e.Source.Repository, Seat: e.Source.Seat}
|
source := buildSource{Repository: e.Source.Repository, Seat: e.Source.Seat}
|
||||||
@@ -286,16 +339,110 @@ func (f following) SourceMoved(ctx context.Context, m link.SourceMoved) error {
|
|||||||
// An empty recorded ref is the repository's default branch, which is what a merge into the base
|
// An empty recorded ref is the repository's default branch, which is what a merge into the base
|
||||||
// branch of the forge's default means.
|
// branch of the forge's default means.
|
||||||
func sourceIs(s inventory.Source, m link.SourceMoved) bool {
|
func sourceIs(s inventory.Source, m link.SourceMoved) bool {
|
||||||
want := strings.ToLower(m.Owner + "/" + m.Repo)
|
if !sameRepository(s.Repository, m) {
|
||||||
repo := strings.ToLower(strings.TrimSuffix(s.Repository, ".git"))
|
|
||||||
matches := repo == want || strings.HasSuffix(repo, "/"+want) ||
|
|
||||||
(m.CloneURL != "" && strings.EqualFold(strings.TrimSuffix(s.Repository, ".git"), strings.TrimSuffix(m.CloneURL, ".git")))
|
|
||||||
if !matches {
|
|
||||||
return false
|
return false
|
||||||
}
|
}
|
||||||
return s.Ref == "" || s.Ref == m.Base
|
return s.Ref == "" || s.Ref == m.Base
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// sameRepository is whether a recorded repository is the one a merge names, in either spelling it
|
||||||
|
// may have been recorded in: a path on the git seat, or the URL it was cloned from.
|
||||||
|
func sameRepository(repository string, m link.SourceMoved) bool {
|
||||||
|
want := strings.ToLower(m.Owner + "/" + m.Repo)
|
||||||
|
repo := strings.ToLower(strings.TrimSuffix(repository, ".git"))
|
||||||
|
return repo == want || strings.HasSuffix(repo, "/"+want) ||
|
||||||
|
(m.CloneURL != "" && repo == strings.ToLower(strings.TrimSuffix(m.CloneURL, ".git")))
|
||||||
|
}
|
||||||
|
|
||||||
|
// readsFrom is whether a module's build read the repository a merge names: the second repository its
|
||||||
|
// recipe packages source from. Its ref must be the branch that moved, or unset — the same rule a
|
||||||
|
// module's own source follows.
|
||||||
|
func readsFrom(read []inventory.ReadRepository, m link.SourceMoved) bool {
|
||||||
|
for _, r := range read {
|
||||||
|
if sameRepository(r.Repository, m) && (r.Ref == "" || r.Ref == m.Base) {
|
||||||
|
return true
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
|
||||||
|
// lastLookAt is the most recent look at this repository by anything built from it.
|
||||||
|
func lastLookAt(entries []inventory.Entry, m link.SourceMoved) time.Time {
|
||||||
|
var newest time.Time
|
||||||
|
for _, e := range entries {
|
||||||
|
if sameRepository(e.Source.Repository, m) && e.Source.Seen.After(newest) {
|
||||||
|
newest = e.Source.Seen
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return newest
|
||||||
|
}
|
||||||
|
|
||||||
|
// whatTheMergeTouched narrows the modules built from a repository to the ones the merge changed.
|
||||||
|
//
|
||||||
|
// **A change inside no module's own directory is a change to what they share.** The forge lists the
|
||||||
|
// files a merge changed; a module is affected when one of them is inside its own directory, when it
|
||||||
|
// is built from the repository's root — everything there is its source — or when some changed file
|
||||||
|
// belongs to no module's directory at all, which is how a shared file, a build recipe or a
|
||||||
|
// dependency at the root rebuilds everything built from that repository.
|
||||||
|
//
|
||||||
|
// A change inside *another* module's directory is that module's business and not this one's, even
|
||||||
|
// when the mesh does not hold that module: `known` is every module this repository is known to hold,
|
||||||
|
// whatever branch it was registered from. That is also the limit of this — a repository whose shared
|
||||||
|
// code sits inside a directory the mesh has never seen a module in reads as shared, and everything
|
||||||
|
// is rebuilt. Rebuilding too much is the safe direction: the fault this whole path exists for is a
|
||||||
|
// mesh that believes it is current and is not (novox/hq 04-ISSUES/131).
|
||||||
|
func whatTheMergeTouched(candidates, known []inventory.Entry, m link.SourceMoved) []inventory.Entry {
|
||||||
|
// Nothing said about the files, or not all of them said: everything built from it is affected.
|
||||||
|
if len(m.Paths) == 0 || m.PathsTruncated {
|
||||||
|
return candidates
|
||||||
|
}
|
||||||
|
var dirs []string
|
||||||
|
for _, e := range known {
|
||||||
|
if e.Source.Path != "" && sameRepository(e.Source.Repository, m) {
|
||||||
|
dirs = append(dirs, e.Source.Path)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
for _, p := range m.Paths {
|
||||||
|
if !insideAny(p, dirs) {
|
||||||
|
return candidates
|
||||||
|
}
|
||||||
|
}
|
||||||
|
var out []inventory.Entry
|
||||||
|
for _, e := range candidates {
|
||||||
|
if e.Source.Path == "" || anyInside(m.Paths, e.Source.Path) {
|
||||||
|
out = append(out, e)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return out
|
||||||
|
}
|
||||||
|
|
||||||
|
// inside is whether a changed file is in a directory: that directory itself, or under it.
|
||||||
|
func inside(path, dir string) bool {
|
||||||
|
dir = strings.Trim(dir, "/")
|
||||||
|
path = strings.TrimPrefix(path, "/")
|
||||||
|
return path == dir || strings.HasPrefix(path, dir+"/")
|
||||||
|
}
|
||||||
|
|
||||||
|
// insideAny is whether a changed file is in any of these directories.
|
||||||
|
func insideAny(path string, dirs []string) bool {
|
||||||
|
for _, dir := range dirs {
|
||||||
|
if inside(path, dir) {
|
||||||
|
return true
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
|
||||||
|
// anyInside is whether any of these changed files is in a directory.
|
||||||
|
func anyInside(paths []string, dir string) bool {
|
||||||
|
for _, p := range paths {
|
||||||
|
if inside(p, dir) {
|
||||||
|
return true
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
|
||||||
// orderByBases is the entries with every base before what stands on it: a module whose build stood
|
// orderByBases is the entries with every base before what stands on it: a module whose build stood
|
||||||
// on another's artifact comes after that module. Entries outside the set are not waited for — they
|
// on another's artifact comes after that module. Entries outside the set are not waited for — they
|
||||||
// are not being rebuilt. Stable for what has no order between it.
|
// are not being rebuilt. Stable for what has no order between it.
|
||||||
@@ -362,3 +509,16 @@ func standsOnModule(e inventory.Entry, module string, against map[string][]strin
|
|||||||
}
|
}
|
||||||
return false
|
return false
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// isHistory is whether a merge made at mergedAt predates the last time the source was seen. A merge
|
||||||
|
// with no time on it is taken as news: refusing it would silence a forge that says less.
|
||||||
|
func isHistory(mergedAt string, seen time.Time) bool {
|
||||||
|
if mergedAt == "" || seen.IsZero() {
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
at, err := time.Parse(time.RFC3339, mergedAt)
|
||||||
|
if err != nil {
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
return at.Before(seen)
|
||||||
|
}
|
||||||
|
|||||||
@@ -66,6 +66,19 @@ func FromEnvironment() (Broker, error) {
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
return Broker{}, err
|
return Broker{}, err
|
||||||
}
|
}
|
||||||
|
// **One bus, one address** (novox/hq ADR 0131). Everything the mesh hands out — a token, a
|
||||||
|
// machine's membership, a person's credential — must name the bus the control plane itself is
|
||||||
|
// connected to; the setting above predates the move and, on a mesh that has moved, still names
|
||||||
|
// the broker it moved from. The first person issued after the move was handed the retired
|
||||||
|
// broker's port and could not connect to anything (2026-09-28).
|
||||||
|
//
|
||||||
|
// Read from the credential rather than from a second setting somebody keeps in step: the
|
||||||
|
// control plane cannot be wrong about where it is connected.
|
||||||
|
if bus, on, err := OnNATS(); err == nil && on {
|
||||||
|
if where := strings.TrimPrefix(BareAddress(bus), "nats://"); where != "" {
|
||||||
|
address = where
|
||||||
|
}
|
||||||
|
}
|
||||||
return Broker{Address: address, Fingerprint: fingerprint}, nil
|
return Broker{Address: address, Fingerprint: fingerprint}, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
+27
-10
@@ -181,6 +181,11 @@ func PermissionsFor(p Principal) (Permissions, error) {
|
|||||||
for _, seat := range meshSeatsTheControllerUses {
|
for _, seat := range meshSeatsTheControllerUses {
|
||||||
pub = append(pub, "mesh.seat."+seat+".accept.>")
|
pub = append(pub, "mesh.seat."+seat+".accept.>")
|
||||||
}
|
}
|
||||||
|
// Every module's tools: **the control plane is the way in** (novox/hq ADR 0095). A person
|
||||||
|
// or an agent asks through it and every question passes one process where an audit
|
||||||
|
// belongs — so it, alone among principals, may call any tool by name. The first `ask` on
|
||||||
|
// the new bus was refused the publish (2026-09-28).
|
||||||
|
pub = append(pub, "mesh.mod.*.tool.>")
|
||||||
|
|
||||||
// The two events it reacts to, and its ack subject on the stream they arrive from
|
// The two events it reacts to, and its ack subject on the stream they arrive from
|
||||||
// (streams.go). **Each named, not a pattern**: `mesh.mod.*.event.>` would make the
|
// (streams.go). **Each named, not a pattern**: `mesh.mod.*.event.>` would make the
|
||||||
@@ -266,9 +271,13 @@ func PermissionsFor(p Principal) (Permissions, error) {
|
|||||||
for _, e := range p.Emits {
|
for _, e := range p.Emits {
|
||||||
pub = append(pub, own+".event."+e)
|
pub = append(pub, own+".event."+e)
|
||||||
}
|
}
|
||||||
for _, t := range p.Serves {
|
// Every tool under its own name, not a list: the tools a module serves are what its code
|
||||||
sub = append(sub, own+".tool."+t)
|
// answers, and a second copy of that list in the manifest would be a second source of
|
||||||
}
|
// truth for the mesh to keep in step (2026-09-28: every module that served a tool was
|
||||||
|
// refused the subscription, because none had written the list twice). Nothing is given
|
||||||
|
// away — no other principal may subscribe this namespace, and a caller's authority is
|
||||||
|
// still granted per tool, by name, on the publish side.
|
||||||
|
sub = append(sub, own+".tool.>")
|
||||||
|
|
||||||
// 2. What it consumes, by the emitter's own subject — an event is addressed to its
|
// 2. What it consumes, by the emitter's own subject — an event is addressed to its
|
||||||
// emitter, because the emitter's identity is the meaning (ADR 0118).
|
// emitter, because the emitter's identity is the meaning (ADR 0118).
|
||||||
@@ -288,11 +297,18 @@ func PermissionsFor(p Principal) (Permissions, error) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// 2c. Its own consumer, which it **pulls**: the runtime asks for the next message and is
|
||||||
|
// answered on its own inbox, so what it needs is to ask about the consumer and to ask it
|
||||||
|
// for messages — its own consumer's name, and no other's. Pulled rather than pushed
|
||||||
|
// because that is the one shape a runtime's client binds without creating anything; the
|
||||||
|
// controller and the hosts are pushed to. Named here rather than through ConsumerFor,
|
||||||
|
// which asks for these permissions to build the consumer and would ask forever. A
|
||||||
|
// subject for a consumer that turns out not to exist grants nothing anybody can use.
|
||||||
|
pub = append(pub,
|
||||||
|
"$JS.API.CONSUMER.INFO."+consumerStream(p)+"."+consumerDurable(p),
|
||||||
|
"$JS.API.CONSUMER.MSG.NEXT."+consumerStream(p)+"."+consumerDurable(p))
|
||||||
|
|
||||||
// 3. Seats it holds: full participation.
|
// 3. Seats it holds: full participation.
|
||||||
// Its consumer's name, not ConsumerFor: that asks for these permissions to build the
|
|
||||||
// consumer, and would ask forever. A subject for a consumer that turns out not to exist
|
|
||||||
// grants nothing anybody can use.
|
|
||||||
sub = append(sub, "_DELIVER."+consumerDurable(p))
|
|
||||||
for _, s := range p.Holds {
|
for _, s := range p.Holds {
|
||||||
// Taking work from the role's queue: the worker consumer it binds (asked about,
|
// Taking work from the role's queue: the worker consumer it binds (asked about,
|
||||||
// delivered on, acknowledged), each on the seat's own stream. The first machine to
|
// delivered on, acknowledged), each on the seat's own stream. The first machine to
|
||||||
@@ -348,9 +364,10 @@ func PermissionsFor(p Principal) (Permissions, error) {
|
|||||||
return Permissions{
|
return Permissions{
|
||||||
Publish: pub,
|
Publish: pub,
|
||||||
Subscribe: sub,
|
Subscribe: sub,
|
||||||
// Only something that serves is ever answering. A pure consumer is granted nothing here.
|
// A module answers what it was asked — a tool call reaches it on its own namespace, so the
|
||||||
AllowResponses: p.Kind == KindModule && (len(p.Serves) > 0 || len(p.Holds) > 0) ||
|
// authority is bounded by having been asked — and so does the controller. A node and a
|
||||||
p.Kind == KindController,
|
// person are never asked anything, and are granted nothing here.
|
||||||
|
AllowResponses: p.Kind == KindModule || p.Kind == KindController,
|
||||||
}, nil
|
}, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -1,6 +1,7 @@
|
|||||||
package broker
|
package broker
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"slices"
|
||||||
"strings"
|
"strings"
|
||||||
"testing"
|
"testing"
|
||||||
)
|
)
|
||||||
@@ -89,17 +90,30 @@ func TestAnInboxIsScopedToItsOwner(t *testing.T) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// A responder answers on the caller's inbox, which it has no permission for. allow_responses is
|
// A responder answers on the caller's inbox, which it has no permission for. allow_responses is
|
||||||
// what makes a scoped inbox workable at all — the authority is bounded by having been asked.
|
// what makes a scoped inbox workable at all — the authority is bounded by having been asked. A
|
||||||
func TestOnlySomethingThatServesMayAnswer(t *testing.T) {
|
// module is asked on its own namespace and may answer; a node and a person are never asked.
|
||||||
serving, _ := PermissionsFor(Principal{Kind: KindModule, Node: "one", Module: "billing",
|
func TestOnlyWhatCanBeAskedMayAnswer(t *testing.T) {
|
||||||
Serves: []string{"status"}, PasswordHash: "x"})
|
module, _ := PermissionsFor(Principal{Kind: KindModule, Node: "one", Module: "audit",
|
||||||
if !serving.AllowResponses {
|
|
||||||
t.Fatal("a module serving a tool cannot answer the caller's inbox")
|
|
||||||
}
|
|
||||||
consumer, _ := PermissionsFor(Principal{Kind: KindModule, Node: "one", Module: "audit",
|
|
||||||
Consumes: []string{"shop.order.placed"}, PasswordHash: "x"})
|
Consumes: []string{"shop.order.placed"}, PasswordHash: "x"})
|
||||||
if consumer.AllowResponses {
|
if !module.AllowResponses {
|
||||||
t.Fatal("a pure consumer was granted the right to answer, which nothing asked it to do")
|
t.Fatal("a module cannot answer a tool call on its own namespace")
|
||||||
|
}
|
||||||
|
node, _ := PermissionsFor(Principal{Kind: KindNode, Node: "one", PasswordHash: "x"})
|
||||||
|
if node.AllowResponses {
|
||||||
|
t.Fatal("a node was granted the right to answer, and nothing asks a node anything")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// A module serves every tool under its own name, and no other module's.
|
||||||
|
func TestAModuleServesItsOwnNamespaceAndNoOthers(t *testing.T) {
|
||||||
|
p, _ := PermissionsFor(Principal{Kind: KindModule, Node: "one", Module: "gitea", PasswordHash: "x"})
|
||||||
|
if !slices.Contains(p.Subscribe, "mesh.mod.gitea.tool.>") {
|
||||||
|
t.Fatalf("a module may not serve its own tools: %v", p.Subscribe)
|
||||||
|
}
|
||||||
|
for _, s := range p.Subscribe {
|
||||||
|
if strings.HasPrefix(s, "mesh.mod.") && !strings.HasPrefix(s, "mesh.mod.gitea.") {
|
||||||
|
t.Fatalf("a module may subscribe another's namespace: %s", s)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -324,3 +338,24 @@ func admits(pattern, subject []string) bool {
|
|||||||
}
|
}
|
||||||
return len(pattern) == len(subject)
|
return len(pattern) == len(subject)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// A module pulls its own consumer — asks about it, asks it for messages — and no other module's.
|
||||||
|
func TestAModulePullsItsOwnConsumerAndNoOthers(t *testing.T) {
|
||||||
|
p, _ := PermissionsFor(Principal{Kind: KindModule, Node: "one", Module: "audit",
|
||||||
|
Consumes: []string{"shop.order.placed"}, PasswordHash: "x"})
|
||||||
|
for _, want := range []string{"$JS.API.CONSUMER.INFO.EVENTS.one_audit", "$JS.API.CONSUMER.MSG.NEXT.EVENTS.one_audit"} {
|
||||||
|
if !slices.Contains(p.Publish, want) {
|
||||||
|
t.Errorf("a module cannot bind its own consumer: %v lacks %s", p.Publish, want)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
for _, s := range p.Publish {
|
||||||
|
if strings.Contains(s, "CONSUMER.") && !strings.HasSuffix(s, ".one_audit") {
|
||||||
|
t.Errorf("a module may reach another consumer: %s", s)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
for _, s := range p.Subscribe {
|
||||||
|
if strings.HasPrefix(s, "_DELIVER.") {
|
||||||
|
t.Errorf("a module is granted a push delivery it never binds: %s", s)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
+9
-7
@@ -24,7 +24,7 @@ accounts {
|
|||||||
jetstream: enabled
|
jetstream: enabled
|
||||||
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.control.>", "mesh.node.>", "mesh.seat.mesh-build-machine.accept.>"] }
|
publish: { allow: ["$JS.ACK.CONTROL.controller.>", "$JS.ACK.EVENTS.controller.>", "$JS.API.>", "_INBOX.enrol.>", "mesh.control.>", "mesh.mod.*.tool.>", "mesh.node.>", "mesh.seat.mesh-build-machine.accept.>"] }
|
||||||
subscribe: { allow: ["$JS.API.>", "_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"] }
|
subscribe: { allow: ["$JS.API.>", "_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"] }
|
||||||
allow_responses: { max: 1, ttl: "1m" }
|
allow_responses: { max: 1, ttl: "1m" }
|
||||||
} }
|
} }
|
||||||
@@ -37,17 +37,19 @@ accounts {
|
|||||||
subscribe: { allow: ["_DELIVER.one", "_INBOX.node.one.>", "mesh.node.one.declare"] }
|
subscribe: { allow: ["_DELIVER.one", "_INBOX.node.one.>", "mesh.node.one.declare"] }
|
||||||
} }
|
} }
|
||||||
{ 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.SEAT_TELEGRAM_SENDER.SEAT_TELEGRAM_SENDER_worker", "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", "mesh.seat.telegram-sender.event.delivered", "mesh.seat.telegram-sender.event.failed"] }
|
||||||
subscribe: { allow: ["_DELIVER.SEAT_TELEGRAM_SENDER_worker", "_DELIVER.one_telegram", "_INBOX.one.telegram.>", "mesh.mod.telegram.tool.status", "mesh.seat.telegram-sender.accept.send"] }
|
subscribe: { allow: ["_DELIVER.SEAT_TELEGRAM_SENDER_worker", "_INBOX.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.>"] }
|
publish: { allow: ["$JS.ACK.EVENTS.two_audit.>", "$JS.API.CONSUMER.INFO.EVENTS.two_audit", "$JS.API.CONSUMER.MSG.NEXT.EVENTS.two_audit"] }
|
||||||
subscribe: { allow: ["_DELIVER.two_audit", "_INBOX.two.audit.>", "mesh.mod.shop.event.order.placed"] }
|
subscribe: { allow: ["_INBOX.two.audit.>", "mesh.mod.audit.tool.>", "mesh.mod.shop.event.order.placed"] }
|
||||||
|
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.>", "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", "mesh.mod.shop.event.order.placed", "mesh.seat.telegram-sender.accept.send"] }
|
||||||
subscribe: { allow: ["_DELIVER.two_shop", "_INBOX.two.shop.>"] }
|
subscribe: { allow: ["_INBOX.two.shop.>", "mesh.mod.shop.tool.>"] }
|
||||||
|
allow_responses: { max: 1, ttl: "1m" }
|
||||||
} }
|
} }
|
||||||
]
|
]
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -61,6 +61,12 @@ type Result struct {
|
|||||||
Commit string
|
Commit string
|
||||||
// Built is each artifact, for reporting.
|
// Built is each artifact, for reporting.
|
||||||
Built []catalogue.Built
|
Built []catalogue.Built
|
||||||
|
|
||||||
|
// Read is every repository this build read source from besides the module's own — the second
|
||||||
|
// repository an artifact's recipe names (ArtifactContext). Reported because the manifest the
|
||||||
|
// mesh keeps carries no build section, so nothing else could say that a merge there is a
|
||||||
|
// change to this module (novox/hq 04-ISSUES/131).
|
||||||
|
Read []catalogue.ArtifactContext
|
||||||
}
|
}
|
||||||
|
|
||||||
// GitCredential is the forge credential a clone may present when the server asks for one.
|
// GitCredential is the forge credential a clone may present when the server asks for one.
|
||||||
@@ -220,7 +226,7 @@ func Build(ctx context.Context, run Runner, publish Publisher,
|
|||||||
}
|
}
|
||||||
say("done", "%s at %s — %d artifact(s) pinned", manifest.Module, short(commit), len(built))
|
say("done", "%s at %s — %d artifact(s) pinned", manifest.Module, short(commit), len(built))
|
||||||
return Result{Manifest: resolved, Commit: commit, Built: built,
|
return Result{Manifest: resolved, Commit: commit, Built: built,
|
||||||
Against: against(within, manifest, stoodOn)}, nil
|
Against: against(within, manifest, stoodOn), Read: readBy(manifest)}, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// Log is where a build says what it is doing, step by step. Nil is silent — the tests pass none,
|
// Log is where a build says what it is doing, step by step. Nil is silent — the tests pass none,
|
||||||
@@ -1016,3 +1022,31 @@ func instructions(recipe string) []string {
|
|||||||
flush()
|
flush()
|
||||||
return out
|
return out
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// readBy is every repository other than the module's own that this build's recipes read source from,
|
||||||
|
// each once and in a fixed order, so two builds of one commit report the same thing the same way.
|
||||||
|
func readBy(manifest catalogue.Manifest) []catalogue.ArtifactContext {
|
||||||
|
if manifest.Build == nil {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
seen := map[string]bool{}
|
||||||
|
var out []catalogue.ArtifactContext
|
||||||
|
for _, a := range manifest.Build.Artifacts {
|
||||||
|
if a.Context == nil || a.Context.Repository == "" {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
key := a.Context.Repository + "#" + a.Context.Ref
|
||||||
|
if seen[key] {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
seen[key] = true
|
||||||
|
out = append(out, *a.Context)
|
||||||
|
}
|
||||||
|
sort.Slice(out, func(i, j int) bool {
|
||||||
|
if out[i].Repository != out[j].Repository {
|
||||||
|
return out[i].Repository < out[j].Repository
|
||||||
|
}
|
||||||
|
return out[i].Ref < out[j].Ref
|
||||||
|
})
|
||||||
|
return out
|
||||||
|
}
|
||||||
|
|||||||
@@ -180,6 +180,20 @@ func (r Registry) MirrorImage(ctx context.Context, from, repository string) (str
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
return "", err
|
return "", err
|
||||||
}
|
}
|
||||||
|
// **Already held is already mirrored.** A base is named by digest, and a digest this registry
|
||||||
|
// holds under the module's repository is the same bytes whatever upstream would say — so
|
||||||
|
// upstream is not asked. Asked every build, the public hub's anonymous pull limit was reached
|
||||||
|
// on the first merge that rebuilt a whole catalogue (2026-09-28), and every module whose base
|
||||||
|
// lives there failed on a copy it did not need.
|
||||||
|
if strings.HasPrefix(where.reference, "sha256:") {
|
||||||
|
held, err := r.has(ctx, "http://"+r.Address+"/v2/"+repository+"/manifests/"+where.reference)
|
||||||
|
if err != nil {
|
||||||
|
return "", fmt.Errorf("asking %s whether it holds %s: %w", r.Address, from, err)
|
||||||
|
}
|
||||||
|
if held {
|
||||||
|
return r.Address + "/" + repository + "@" + where.reference, nil
|
||||||
|
}
|
||||||
|
}
|
||||||
src := &source{client: r.client()}
|
src := &source{client: r.client()}
|
||||||
digest, err := r.copyManifest(ctx, src, where, where.reference, repository)
|
digest, err := r.copyManifest(ctx, src, where, where.reference, repository)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
|||||||
@@ -103,6 +103,12 @@ func (m *theMeshsRegistry) handler() http.Handler {
|
|||||||
m.mu.Lock()
|
m.mu.Lock()
|
||||||
defer m.mu.Unlock()
|
defer m.mu.Unlock()
|
||||||
switch {
|
switch {
|
||||||
|
case r.Method == http.MethodHead && strings.Contains(r.URL.Path, "/manifests/"):
|
||||||
|
if _, ok := m.manifests[r.URL.Path[strings.LastIndex(r.URL.Path, "/")+1:]]; ok {
|
||||||
|
w.WriteHeader(http.StatusOK)
|
||||||
|
} else {
|
||||||
|
w.WriteHeader(http.StatusNotFound)
|
||||||
|
}
|
||||||
case r.Method == http.MethodHead && strings.Contains(r.URL.Path, "/blobs/"):
|
case r.Method == http.MethodHead && strings.Contains(r.URL.Path, "/blobs/"):
|
||||||
if _, ok := m.blobs[r.URL.Path[strings.LastIndex(r.URL.Path, "/")+1:]]; ok {
|
if _, ok := m.blobs[r.URL.Path[strings.LastIndex(r.URL.Path, "/")+1:]]; ok {
|
||||||
w.WriteHeader(http.StatusOK)
|
w.WriteHeader(http.StatusOK)
|
||||||
@@ -219,3 +225,27 @@ func TestATagBeforeTheDigestIsNotPartOfTheRepository(t *testing.T) {
|
|||||||
t.Fatalf("got %+v", got)
|
t.Fatalf("got %+v", got)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// A base this registry already holds by digest is not asked of upstream at all: the public hub
|
||||||
|
// limits anonymous pulls, and a catalogue rebuilt on one merge asked it once per module.
|
||||||
|
func TestABaseAlreadyHeldIsNotAskedOfUpstream(t *testing.T) {
|
||||||
|
src, indexDigest, _ := anUpstreamRegistry(t)
|
||||||
|
dst := &theMeshsRegistry{blobs: map[string][]byte{}, manifests: map[string][]byte{}}
|
||||||
|
dstServer := httptest.NewServer(dst.handler())
|
||||||
|
defer dstServer.Close()
|
||||||
|
address := strings.TrimPrefix(dstServer.URL, "http://")
|
||||||
|
r := Registry{Address: address, HTTP: src.Client()}
|
||||||
|
host := strings.TrimPrefix(src.URL, "http://")
|
||||||
|
if _, err := r.MirrorImage(context.Background(), host+"/library/thing:latest", "hello-web/server"); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
// Upstream gone: the pinned base is answered from what the mesh holds.
|
||||||
|
src.Close()
|
||||||
|
reference, err := r.MirrorImage(context.Background(), host+"/library/thing@"+indexDigest, "hello-web/server")
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("a base the registry holds was asked of an upstream that is gone: %v", err)
|
||||||
|
}
|
||||||
|
if reference != address+"/hello-web/server@"+indexDigest {
|
||||||
|
t.Fatalf("pinned as %q", reference)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -179,3 +179,24 @@ func TestTheBasesABuildWasHandedAreWhatItStoodOn(t *testing.T) {
|
|||||||
t.Fatalf("the bases the build was handed were not what it stood on: %v", got)
|
t.Fatalf("the bases the build was handed were not what it stood on: %v", got)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// What a build read besides its module's own repository is the second repository its recipes name,
|
||||||
|
// each once: a module that packages source living elsewhere is affected when that source moves.
|
||||||
|
func TestWhatABuildReadIsTheRepositoriesItsRecipesName(t *testing.T) {
|
||||||
|
elsewhere := catalogue.ArtifactContext{Repository: "http://forge.internal:20000/novox/mesh-controller.git", Ref: "main"}
|
||||||
|
manifest := catalogue.Manifest{
|
||||||
|
Module: "builder",
|
||||||
|
Build: &catalogue.Build{Artifacts: []catalogue.Artifact{
|
||||||
|
{Name: "server", Kind: catalogue.ArtifactImage, From: "Dockerfile", Context: &elsewhere},
|
||||||
|
{Name: "tools", Kind: catalogue.ArtifactImage, From: "Dockerfile", Context: &elsewhere},
|
||||||
|
{Name: "config", Kind: catalogue.ArtifactArchive, From: "etc"},
|
||||||
|
}},
|
||||||
|
}
|
||||||
|
read := readBy(manifest)
|
||||||
|
if len(read) != 1 || read[0] != elsewhere {
|
||||||
|
t.Fatalf("the repositories this build read are %+v", read)
|
||||||
|
}
|
||||||
|
if readBy(catalogue.Manifest{Module: "gitea", Build: &catalogue.Build{}}) != nil {
|
||||||
|
t.Fatal("a module whose recipes name no other repository read one")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -36,12 +36,21 @@ type Build struct {
|
|||||||
// Against is every artifact this build stood on, as references rather than module names —
|
// Against is every artifact this build stood on, as references rather than module names —
|
||||||
// what makes a build edge derived rather than declared (ADR 0009).
|
// what makes a build edge derived rather than declared (ADR 0009).
|
||||||
Against []string
|
Against []string
|
||||||
|
// Read is every repository this build read source from besides the module's own (novox/hq
|
||||||
|
// 04-ISSUES/131), at the ref it read.
|
||||||
|
Read []ReadRepository
|
||||||
// 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
|
||||||
At time.Time
|
At time.Time
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// ReadRepository is a repository a build read source from besides the module's own.
|
||||||
|
type ReadRepository struct {
|
||||||
|
Repository string `json:"repository"`
|
||||||
|
Ref string `json:"ref,omitempty"`
|
||||||
|
}
|
||||||
|
|
||||||
// Artifact is one thing a build published.
|
// Artifact is one thing a build published.
|
||||||
type Artifact struct {
|
type Artifact struct {
|
||||||
Name string `json:"name"`
|
Name string `json:"name"`
|
||||||
@@ -66,17 +75,21 @@ func (i *Inventory) RecordBuild(ctx context.Context, b Build) error {
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
read, err := json.Marshal(b.Read)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
var module *string
|
var module *string
|
||||||
if b.Module != "" {
|
if b.Module != "" {
|
||||||
module = &b.Module
|
module = &b.Module
|
||||||
}
|
}
|
||||||
_, 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)
|
source_path, manifest, built_against, built_contexts)
|
||||||
values ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11)
|
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)
|
b.Path, manifestOrNil(b.Manifest), against, read)
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -199,6 +212,44 @@ func (i *Inventory) BuiltAgainst(ctx context.Context) (map[string][]string, erro
|
|||||||
return against, rows.Err()
|
return against, rows.Err()
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// ReadRepositories is what each module's newest successful build read source from besides its own
|
||||||
|
// repository, by module name.
|
||||||
|
//
|
||||||
|
// The mirror of BuiltAgainst, and derived the same way and for the same reason: a merge into a
|
||||||
|
// repository a module only packages is a change to that module, and the manifest the mesh keeps
|
||||||
|
// carries nothing that would say so (novox/hq 04-ISSUES/131).
|
||||||
|
func (i *Inventory) ReadRepositories(ctx context.Context) (map[string][]ReadRepository, error) {
|
||||||
|
rows, err := i.store.Pool().Query(ctx,
|
||||||
|
`select distinct on (module) module, built_contexts
|
||||||
|
from build
|
||||||
|
where module is not null and module <> '' and failed = ''
|
||||||
|
order by module, at desc`)
|
||||||
|
if err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
defer rows.Close()
|
||||||
|
|
||||||
|
read := map[string][]ReadRepository{}
|
||||||
|
for rows.Next() {
|
||||||
|
var module string
|
||||||
|
var raw []byte
|
||||||
|
if err := rows.Scan(&module, &raw); err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
if len(raw) == 0 {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
var of []ReadRepository
|
||||||
|
if err := json.Unmarshal(raw, &of); err != nil {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
if len(of) > 0 {
|
||||||
|
read[module] = of
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return read, rows.Err()
|
||||||
|
}
|
||||||
|
|
||||||
// manifestOrNil keeps the difference between "declared nothing" and "predates this being kept".
|
// manifestOrNil keeps the difference between "declared nothing" and "predates this being kept".
|
||||||
//
|
//
|
||||||
// A build recorded before the mesh kept manifests has no manifest, and that is not the same as one
|
// A build recorded before the mesh kept manifests has no manifest, and that is not the same as one
|
||||||
|
|||||||
@@ -136,3 +136,36 @@ func ids(builds []Build) []string {
|
|||||||
}
|
}
|
||||||
return out
|
return out
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// What a build read besides its module's own repository comes back for the newest build of each
|
||||||
|
// module, and only for builds that worked. Nothing recorded is absent rather than empty, which is how
|
||||||
|
// a build made before the mesh kept this is told from one that read nothing (novox/hq 04-ISSUES/131).
|
||||||
|
func TestWhatABuildReadComesBackForTheNewestBuildOfEachModule(t *testing.T) {
|
||||||
|
inv := fresh(t)
|
||||||
|
ctx := context.Background()
|
||||||
|
older := aBuild("older", "builder", "")
|
||||||
|
older.Read = []ReadRepository{{Repository: "novox/mesh-controller", Ref: "release"}}
|
||||||
|
newer := aBuild("newer", "builder", "")
|
||||||
|
newer.Read = []ReadRepository{{Repository: "novox/mesh-controller", Ref: "main"}}
|
||||||
|
plain := aBuild("plain", "gitea", "")
|
||||||
|
failed := aBuild("failed", "route-proxy", "cannot clone")
|
||||||
|
failed.Read = []ReadRepository{{Repository: "novox/mesh-controller", Ref: "main"}}
|
||||||
|
for _, b := range []Build{older, newer, plain, failed} {
|
||||||
|
if err := inv.RecordBuild(ctx, b); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
read, err := inv.ReadRepositories(ctx)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if len(read["builder"]) != 1 || read["builder"][0].Ref != "main" {
|
||||||
|
t.Fatalf("the newest build's reading is %+v", read["builder"])
|
||||||
|
}
|
||||||
|
if _, has := read["gitea"]; has {
|
||||||
|
t.Fatalf("a build that read nothing but its own repository reads as %+v", read["gitea"])
|
||||||
|
}
|
||||||
|
if _, has := read["route-proxy"]; has {
|
||||||
|
t.Fatal("a failed build's reading was kept as what that module reads")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -86,9 +86,11 @@ func TestAnAssignedModuleBecomesAUserWithWhatItDeclared(t *testing.T) {
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
|
// Its tools are every one under its own name — the list in the manifest is a person's
|
||||||
|
// vocabulary for asking, not the module's permission to answer.
|
||||||
if !granted(perms.Publish, "mesh.mod.shop.event.order.placed") ||
|
if !granted(perms.Publish, "mesh.mod.shop.event.order.placed") ||
|
||||||
!granted(perms.Publish, "mesh.seat.telegram-sender.accept.send") ||
|
!granted(perms.Publish, "mesh.seat.telegram-sender.accept.send") ||
|
||||||
!granted(perms.Subscribe, "mesh.mod.shop.tool.price") {
|
!granted(perms.Subscribe, "mesh.mod.shop.tool.>") {
|
||||||
t.Fatalf("one.shop's authority is not what it declared: %+v", perms)
|
t.Fatalf("one.shop's authority is not what it declared: %+v", perms)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -8,6 +8,7 @@ import (
|
|||||||
"sort"
|
"sort"
|
||||||
"strconv"
|
"strconv"
|
||||||
"strings"
|
"strings"
|
||||||
|
"time"
|
||||||
|
|
||||||
"github.com/jackc/pgx/v5"
|
"github.com/jackc/pgx/v5"
|
||||||
"github.com/novox/mesh-controller/internal/catalogue"
|
"github.com/novox/mesh-controller/internal/catalogue"
|
||||||
@@ -38,6 +39,9 @@ type Source struct {
|
|||||||
BuiltFrom string
|
BuiltFrom string
|
||||||
// Head is the newest commit the source is known to have.
|
// Head is the newest commit the source is known to have.
|
||||||
Head string
|
Head string
|
||||||
|
// Seen is when the source was last looked at — by a build, by hand, or by the forge saying it
|
||||||
|
// moved. What a late report of an older move is judged against.
|
||||||
|
Seen 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.
|
||||||
@@ -936,12 +940,13 @@ func (i *Inventory) Catalogued(ctx context.Context) ([]Entry, error) {
|
|||||||
`select m.name, m.manifest,
|
`select m.name, m.manifest,
|
||||||
coalesce(m.source, ''), m.source_path, m.source_seat, coalesce(m.ref, ''),
|
coalesce(m.source, ''), m.source_path, m.source_seat, coalesce(m.ref, ''),
|
||||||
coalesce(m.built_from, ''), coalesce(m.source_head, ''),
|
coalesce(m.built_from, ''), coalesce(m.source_head, ''),
|
||||||
|
coalesce(m.source_seen, to_timestamp(0)),
|
||||||
coalesce(array_agg(n.name order by n.name) filter (where n.name is not null), '{}')
|
coalesce(array_agg(n.name order by n.name) filter (where n.name is not null), '{}')
|
||||||
from module m
|
from module m
|
||||||
left join assignment a on a.module = m.name
|
left join assignment a on a.module = m.name
|
||||||
left join node n on n.id = a.node
|
left join node n on n.id = a.node
|
||||||
group by m.name, m.manifest, m.source, m.source_path, m.source_seat, m.ref, m.built_from,
|
group by m.name, m.manifest, m.source, m.source_path, m.source_seat, m.ref, m.built_from,
|
||||||
m.source_head
|
m.source_head, m.source_seen
|
||||||
order by m.name`)
|
order by m.name`)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
@@ -955,9 +960,12 @@ func (i *Inventory) Catalogued(ctx context.Context) ([]Entry, error) {
|
|||||||
var source Source
|
var source Source
|
||||||
var on []string
|
var on []string
|
||||||
if err := rows.Scan(&name, &raw, &source.Repository, &source.Path, &source.Seat, &source.Ref,
|
if err := rows.Scan(&name, &raw, &source.Repository, &source.Path, &source.Seat, &source.Ref,
|
||||||
&source.BuiltFrom, &source.Head, &on); err != nil {
|
&source.BuiltFrom, &source.Head, &source.Seen, &on); err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
|
if source.Seen.Unix() == 0 {
|
||||||
|
source.Seen = time.Time{}
|
||||||
|
}
|
||||||
var m catalogue.Manifest
|
var m catalogue.Manifest
|
||||||
if err := json.Unmarshal(raw, &m); err != nil {
|
if err := json.Unmarshal(raw, &m); err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
|
|||||||
@@ -0,0 +1,13 @@
|
|||||||
|
-- A build says which repositories it read, so a merge can find everything it affects.
|
||||||
|
--
|
||||||
|
-- A module is built from the one repository the mesh records — where its module.json lives — and some
|
||||||
|
-- modules' recipes reach into a second for the source they package: the packaging and the source are
|
||||||
|
-- allowed to live apart (catalogue's ArtifactContext). That second repository is named in the manifest
|
||||||
|
-- the build read, and the manifest the mesh *keeps* carries no build section, so nothing on the mesh
|
||||||
|
-- could say that a merge into the other repository is a change to this module at all. Two modules are
|
||||||
|
-- built from the control plane's own repository, and neither had ever been rebuilt when it moved
|
||||||
|
-- (novox/hq 04-ISSUES/131).
|
||||||
|
--
|
||||||
|
-- Nullable, like the two derived columns beside it: null is a build recorded before the mesh kept
|
||||||
|
-- this, which is not the same as a build that read nothing but its module's own repository.
|
||||||
|
alter table build add column built_contexts jsonb;
|
||||||
@@ -88,10 +88,23 @@ type BuildResult struct {
|
|||||||
// (novox/hq ADR 0009). The catalogue turns these into edges; nothing else need care.
|
// (novox/hq ADR 0009). The catalogue turns these into edges; nothing else need care.
|
||||||
Against []string `json:"against,omitempty"`
|
Against []string `json:"against,omitempty"`
|
||||||
|
|
||||||
|
// Read is every repository this build read source from besides the module's own. A module whose
|
||||||
|
// recipe packages source that lives elsewhere is affected when that repository moves, and the
|
||||||
|
// manifest the mesh keeps says nothing about it (novox/hq 04-ISSUES/131).
|
||||||
|
Read []ReadRepository `json:"read,omitempty"`
|
||||||
|
|
||||||
// Failed is why, when it did.
|
// Failed is why, when it did.
|
||||||
Failed string `json:"failed,omitempty"`
|
Failed string `json:"failed,omitempty"`
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// ReadRepository is a repository a build read source from besides the module's own, at the branch,
|
||||||
|
// tag or commit it read. Spelled here as well as in the catalogue and the inventory, for the reason
|
||||||
|
// MadeArtifact is: one direction of dependency.
|
||||||
|
type ReadRepository struct {
|
||||||
|
Repository string `json:"repository"`
|
||||||
|
Ref string `json:"ref,omitempty"`
|
||||||
|
}
|
||||||
|
|
||||||
// MadeArtifact is one thing a build produced, as a person would want it reported.
|
// MadeArtifact is one thing a build produced, as a person would want it reported.
|
||||||
type MadeArtifact struct {
|
type MadeArtifact struct {
|
||||||
Name string `json:"name"`
|
Name string `json:"name"`
|
||||||
|
|||||||
@@ -128,6 +128,18 @@ type SourceMoved struct {
|
|||||||
Commit string `json:"merge_commit_sha"`
|
Commit string `json:"merge_commit_sha"`
|
||||||
CloneURL string `json:"clone_url"`
|
CloneURL string `json:"clone_url"`
|
||||||
HTMLURL string `json:"html_url"`
|
HTMLURL string `json:"html_url"`
|
||||||
|
// MergedAt is when the forge merged it, RFC 3339. What decides whether this is news.
|
||||||
|
MergedAt string `json:"merged_at"`
|
||||||
|
|
||||||
|
// Paths are the files the merge changed, from the repository's root. Empty means the forge said
|
||||||
|
// nothing about them, and every module built from the repository is treated as affected.
|
||||||
|
Paths []string `json:"paths,omitempty"`
|
||||||
|
|
||||||
|
// PathsTruncated says the merge changed more files than the forge was asked to list, so Paths is
|
||||||
|
// a beginning rather than the whole change — and again, everything is treated as affected. Said
|
||||||
|
// rather than inferred from a round number, because "this is all of it" and "this is as much as
|
||||||
|
// I asked for" are the difference between rebuilding a module and leaving it stale.
|
||||||
|
PathsTruncated bool `json:"paths_truncated,omitempty"`
|
||||||
}
|
}
|
||||||
|
|
||||||
type Upgraded struct {
|
type Upgraded struct {
|
||||||
|
|||||||
@@ -154,10 +154,47 @@ func (n *natsInbound) deliver(ctx context.Context, act func(context.Context, Con
|
|||||||
m.seq = meta.Sequence.Stream
|
m.seq = meta.Sequence.Stream
|
||||||
m.delivered = meta.NumDelivered
|
m.delivered = meta.NumDelivered
|
||||||
}
|
}
|
||||||
|
// **Work that outlives the acknowledgement window says so while it runs.**
|
||||||
|
//
|
||||||
|
// The bus waits a fixed time to be told a message was taken, and then hands it to whoever
|
||||||
|
// consumes next — which is right for a consumer that died and wrong for one that is busy.
|
||||||
|
// Acting on a merge builds every module the merge changed: minutes of work against a
|
||||||
|
// thirty-second window. So the same merge was handed over again while the first build was
|
||||||
|
// still running, and again after that — on 2026-09-28 one merge ran the mesh's whole
|
||||||
|
// catalogue five times over and exhausted a public registry's pull limit.
|
||||||
|
//
|
||||||
|
// Here rather than in each handler, because the window belongs to the transport and every
|
||||||
|
// handler would otherwise have to remember it. It changes nothing about a handler that
|
||||||
|
// dies: a message is kept alive only while this goroutine is, so a controller that stops
|
||||||
|
// stops saying so, and the bus redelivers exactly as it should.
|
||||||
|
working := make(chan struct{})
|
||||||
|
defer close(working)
|
||||||
|
go stillWorking(msg, working)
|
||||||
}
|
}
|
||||||
act(ctx, m)
|
act(ctx, m)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// heartbeatWhileWorking is how often a handler still running tells the bus so — comfortably inside
|
||||||
|
// the shortest acknowledgement window the mesh gives any of its consumers.
|
||||||
|
const heartbeatWhileWorking = 10 * time.Second
|
||||||
|
|
||||||
|
// stillWorking keeps one message alive until the work on it returns.
|
||||||
|
//
|
||||||
|
// An error is not worth reporting: what the bus does when it is not told is redeliver, which is
|
||||||
|
// exactly what happens if this fails, and the handler's own outcome is the thing worth logging.
|
||||||
|
func stillWorking(msg *nats.Msg, done <-chan struct{}) {
|
||||||
|
tick := time.NewTicker(heartbeatWhileWorking)
|
||||||
|
defer tick.Stop()
|
||||||
|
for {
|
||||||
|
select {
|
||||||
|
case <-done:
|
||||||
|
return
|
||||||
|
case <-tick.C:
|
||||||
|
_ = msg.InProgress()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
// kindOfSubject is how this transport's addressing becomes what the mesh calls a message.
|
// kindOfSubject is how this transport's addressing becomes what the mesh calls a message.
|
||||||
//
|
//
|
||||||
// By subject, which is the only thing the server enforces: a body claiming to be a report does not
|
// By subject, which is the only thing the server enforces: a body claiming to be a report does not
|
||||||
|
|||||||
@@ -395,3 +395,60 @@ func TestEverySubjectTheControllerFollowsDecodesToAKind(t *testing.T) {
|
|||||||
t.Error("a subject nobody follows decoded to a kind")
|
t.Error("a subject nobody follows decoded to a kind")
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// **A handler slower than the acknowledgement window is not handed its message again.**
|
||||||
|
//
|
||||||
|
// The bus waits a fixed time to be told a message was taken and then redelivers, which is right for
|
||||||
|
// a consumer that died and wrong for one that is busy. Acting on a merge builds modules — minutes
|
||||||
|
// against a thirty-second window — and the same merge was handed over five times while the first
|
||||||
|
// build was still running (2026-09-28). Here the window is two seconds and the work takes six.
|
||||||
|
func TestNatsWorkSlowerThanTheWindowIsNotHandedOverAgain(t *testing.T) {
|
||||||
|
js := aBus(t)
|
||||||
|
// The controller's own consumer, with a window short enough to outlive in a test.
|
||||||
|
if err := js.EnsureConsumer(broker.Consumer{
|
||||||
|
Name: broker.ControllerName, Stream: "CONTROL", Push: true, AckWaitSeconds: 2,
|
||||||
|
Why: "a window short enough to outlive in a test",
|
||||||
|
}); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
var mu sync.Mutex
|
||||||
|
handled := 0
|
||||||
|
slow := make(chan struct{})
|
||||||
|
s, stop := servingOn(t, js, nil)
|
||||||
|
defer stop()
|
||||||
|
s.listener = slowly{func() {
|
||||||
|
mu.Lock()
|
||||||
|
handled++
|
||||||
|
first := handled == 1
|
||||||
|
mu.Unlock()
|
||||||
|
if first {
|
||||||
|
time.Sleep(6 * time.Second)
|
||||||
|
close(slow)
|
||||||
|
}
|
||||||
|
}}
|
||||||
|
|
||||||
|
body, err := json.Marshal(Report{Node: "anchor", Declared: "d1", Applied: []string{"store"}})
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if _, err := js.Context().Publish(ReportSubject("anchor"), body); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
select {
|
||||||
|
case <-slow:
|
||||||
|
case <-time.After(30 * time.Second):
|
||||||
|
t.Fatal("the slow work never finished")
|
||||||
|
}
|
||||||
|
// A moment for a redelivery to arrive, if the bus were going to send one.
|
||||||
|
time.Sleep(3 * time.Second)
|
||||||
|
mu.Lock()
|
||||||
|
defer mu.Unlock()
|
||||||
|
if handled != 1 {
|
||||||
|
t.Fatalf("one report was handled %d times, so slow work is run again while it is running", handled)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// slowly is a listener that runs whatever it was given.
|
||||||
|
type slowly struct{ work func() }
|
||||||
|
|
||||||
|
func (s slowly) Heard(context.Context, Report) error { s.work(); return nil }
|
||||||
|
|||||||
Reference in New Issue
Block a user