Compare commits

..
Author SHA1 Message Date
mesh-admin cec792ce9d Merge pull request 'A manifest HEAD says what it accepts, or the registry answers 404' (#127) from fix/a-manifest-head-says-what-it-accepts into main 2026-09-28 11:02:10 +00:00
jschoubben 338d033632 A manifest HEAD says what it accepts, or the registry answers 404
The check that skips copying a base the mesh already holds asked with no Accept header, and a
registry answers a manifest only in a media type the caller named: the same digest answered 200 with
the manifest types and 404 without them. So the builder concluded it held nothing, copied every
vendor base again, and exhausted the public hub's pull limit a second time today.

The test could not have caught it, because the fake registry answered a manifest HEAD regardless of
Accept — more permissive than the thing it stands in for. It is now as strict as a real registry, and
fails without the fix.
2026-09-28 13:02:08 +02:00
mesh-admin 1be926cec4 Merge pull request 'The control plane prepares its own state, like any module' (#126) from feat/the-control-plane-prepares-its-own-state into main 2026-09-28 10:51:30 +00:00
jschoubben 2134768dfe The control plane prepares its own state, like any module
Now that every parser on the mesh knows the word, the control plane's manifest says it. Its schema
stops being a special case: the mesh derives the step from its own resource and gates its server on
it, which is the failure of novox/hq 04-ISSUES/133 closed by the mechanism rather than by a
hand-written step in one manifest.
2026-09-28 12:51:27 +02:00
mesh-admin ef825688ee Merge pull request 'The control plane learns 'prepares' one release before its manifest uses it' (#125) from fix/the-word-ships-before-the-manifest-uses-it into main 2026-09-28 10:49:36 +00:00
jschoubben 77a14360df The control plane learns 'prepares' one release before its manifest uses it
A manifest word has to reach every parser before a manifest carries it. The builder refused
mesh-controller's manifest with `unknown field "prepares"` until it was rebuilt; then the running
control plane could not read the manifest in the build result either, so the build was recorded with
no module and the version never moved. Strict parsing is deliberate (novox/hq 04-ISSUES/003), so the
word ships first and a manifest uses it next: this takes `prepares` back out of the control plane's
own manifest, leaving the code that understands it, and the manifest says it again once this is
running everywhere.
2026-09-28 12:49:33 +02:00
mesh-admin 9be2fb4750 Merge pull request 'A version prepares its state before it runs' (#124) from feat/a-version-prepares-its-state into main 2026-09-28 10:43:17 +00:00
jschoubben 5a963aec10 A version prepares its state before it runs
The mesh derives the preparation from the module's own resource instead of each module hand-writing
a step beside it (novox/hq ADR 0135). A manifest says one word — `prepares` — and the mesh runs that
module's own program in its preparation mode, in the module's own context: the same image, the same
environment, the same mounts, because it is the same code. A published port and a fixed address are
taken away rather than copied, since the version being replaced still holds them.

One word for every kind of module: a Go binary receives `prepare` as its argument, a bundle receives
it through the runtime whose entry takes the same word. The control plane answers it like anything
else — its own schema stops being a special case, and its hand-written step is gone.
2026-09-28 12:43:15 +02:00
mesh-admin 45d1c28a28 Merge pull request 'The control plane migrates before it serves' (#123) from fix/the-control-plane-migrates-before-it-serves into main 2026-09-28 08:27:44 +00:00
jschoubben a3e7683c63 The control plane migrates before it serves
The mesh replaced its own control plane with a build carrying a migration, applied none of it, and
then refused every build it recorded for three quarters of an hour while reporting itself healthy
(novox/hq 04-ISSUES/133). The module now declares the step ADR 0052 prescribes: a run-once
`migrate` before the server, re-run whenever the image moves because the image is part of a step's
digest, and gating — a migration that fails stops the new server from starting rather than letting
it serve against a schema it does not have.
2026-09-28 10:27:42 +02:00
mesh-admin da394b45e6 Merge pull request 'Work slower than the window says so, and one address is the bus's' (#122) from fix/work-longer-than-the-window-says-so into main 2026-09-28 07:51:52 +00:00
jschoubben 2f3bfda8c0 Work slower than the window says so, and one address is the bus's
Three faults the mesh's own logs showed this morning. A handler that outlives the acknowledgement
window was handed its message again while it was still working: acting on a merge builds modules,
minutes against a thirty-second window, so one merge ran the whole catalogue five times over. The
transport now says the work is in progress while it runs, which is where the window belongs.

Everything the mesh hands out — a token, a membership, a person's credential — took its address
from the enrolment setting, which 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. There is one bus, and
its address is the one the control plane is connected to.

And `operator issue` documented an argument order its parser refused.
2026-09-28 09:51:50 +02:00
mesh-admin 208388978a Merge pull request 'A merge rebuilds what it changed, and what packages it' (#121) from feat/a-merge-rebuilds-what-it-changed into main 2026-09-28 07:20:03 +00:00
13 changed files with 395 additions and 6 deletions
+5 -1
View File
@@ -76,7 +76,10 @@ func run() error {
return pinCommand(ctx, args[1:], true)
case "unpin":
return pinCommand(ctx, args[1:], false)
case "migrate":
// `prepare` is how the mesh asks any module to bring its state to the shape this version needs
// (novox/hq ADR 0135), and the control plane answers it the same way as everything else — its
// own schema is not a special case. `migrate` remains the word a person types.
case "prepare", "migrate":
return migrate(ctx)
case "node":
return nodeCommand(ctx, args[1:])
@@ -138,6 +141,7 @@ func usage() {
fmt.Fprint(os.Stderr, `mesh-controller — the control plane
migrate bring each context's schema up to date
prepare the same, asked the way the mesh asks any module (ADR 0135)
node add <name> [--adopted] create a node record; --adopted: the machine is in use
node list the nodes this mesh knows about
node show <name> what one machine reported it can do, and why
+7 -3
View File
@@ -189,13 +189,17 @@ func readPrivateKey(path string) (string, error) {
func personIssue(ctx context.Context, args []string) error {
set := flag.NewFlagSet("operator issue", flag.ContinueOnError)
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
}
if set.NArg() != 1 {
if len(positionals) != 1 {
return errors.New("operator issue <name> --invokes <tool,tool|*>")
}
name := set.Arg(0)
name := positionals[0]
if *invokes == "" {
return errors.New(
"say what this person may call: --invokes mesh-catalog.catalog_tools,gitea.repo_create, " +
+9
View File
@@ -245,10 +245,12 @@ func (f following) SourceMoved(ctx context.Context, m link.SourceMoved) error {
// 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 {
switch {
case sourceIs(e.Source, m):
if e.Source.BuiltFrom == m.Commit {
already++
continue
}
// **A merge older than the last look at the source is history, not a move.** The forge
@@ -264,6 +266,13 @@ func (f following) SourceMoved(ctx context.Context, m link.SourceMoved) error {
}
}
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
+13
View File
@@ -66,6 +66,19 @@ func FromEnvironment() (Broker, error) {
if err != nil {
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
}
+1 -1
View File
@@ -186,7 +186,7 @@ func (r Registry) MirrorImage(ctx context.Context, from, repository string) (str
// 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)
held, err := r.has(ctx, "http://"+r.Address+"/v2/"+repository+"/manifests/"+where.reference, manifestAccept)
if err != nil {
return "", fmt.Errorf("asking %s whether it holds %s: %w", r.Address, from, err)
}
+8
View File
@@ -104,6 +104,14 @@ func (m *theMeshsRegistry) handler() http.Handler {
defer m.mu.Unlock()
switch {
case r.Method == http.MethodHead && strings.Contains(r.URL.Path, "/manifests/"):
// **As strictly as a real registry.** A manifest is answered only in a media type the
// caller named; a request with no Accept is answered as if nothing were there. The fake
// used to answer regardless, which is why it could not catch a check that asked without
// one — and the mesh copied every base again (2026-09-28).
if !strings.Contains(r.Header.Get("Accept"), "manifest") && !strings.Contains(r.Header.Get("Accept"), "index") {
w.WriteHeader(http.StatusNotFound)
return
}
if _, ok := m.manifests[r.URL.Path[strings.LastIndex(r.URL.Path, "/")+1:]]; ok {
w.WriteHeader(http.StatusOK)
} else {
+13 -1
View File
@@ -122,11 +122,23 @@ func (r Registry) PublishArchive(ctx context.Context, repository string, body []
return final, nil
}
func (r Registry) has(ctx context.Context, url string) (bool, error) {
// has is whether this registry already holds what is at that URL.
//
// **A manifest HEAD must say what it accepts.** A registry answers a manifest request only in a media
// type the caller named, and a bare HEAD — no Accept at all — is answered 404 for a manifest it holds
// perfectly well. Measured against the mesh's own registry (2026-09-28): the same digest answered 200
// with the manifest media types and 404 without them, so a check written without them concluded the
// registry held nothing, copied every base again, and exhausted the public hub's pull limit. A blob
// needs no Accept, which is why this went unnoticed: the same helper was right for blobs and wrong
// for manifests.
func (r Registry) has(ctx context.Context, url string, accept ...string) (bool, error) {
request, err := http.NewRequestWithContext(ctx, http.MethodHead, url, nil)
if err != nil {
return false, err
}
for _, media := range accept {
request.Header.Set("Accept", media)
}
response, err := r.client().Do(request)
if err != nil {
return false, fmt.Errorf("cannot reach the registry at %s: %w", r.Address, err)
+85
View File
@@ -623,6 +623,9 @@ func (r Resolution) compose(with Rendering, owner map[string]string) ([]map[stri
// one of them as its environment without saying so (ADR 0086, issue 041).
secretFiles := secretFilesOf(resources)
// Which of this module's resources its preparation runs before, if it prepares anything.
prepareBefore := preparationTarget(m)
for _, unsettled := range resources {
resource, err := ApplySettings(unsettled, with.Settings[m.Module])
if err != nil {
@@ -708,6 +711,18 @@ func (r Resolution) compose(with Rendering, owner map[string]string) ([]map[stri
if renamed := reflectsRenamed(m.Module, resource["reload-on"]); renamed != nil {
copied["reload-on"] = renamed
}
// **A version prepares its state before it runs** (novox/hq ADR 0135). Derived from the
// module's own resource rather than declared beside it: what prepares the state is the
// module's own code, so what it is given has to be what that code is given — and a
// second resource written by hand is a second copy to drift from the first. Placed
// immediately before it, because a run-once step stops everything the declaration
// places after it (ADR 0052), which is how a version whose preparation failed does not
// serve.
if prepareBefore != "" && fmt.Sprint(resource["id"]) == prepareBefore {
step := prepared(copied)
owner[fmt.Sprint(step["id"])] = m.Module
out = append(out, step)
}
owner[fmt.Sprint(copied["id"])] = m.Module
out = append(out, copied)
}
@@ -1610,3 +1625,73 @@ func atMachinePort(serves map[string]any, module string, ports map[string]map[in
func AtPublishedPort(values map[string]any, module string, published map[int]int) map[string]any {
return atMachinePort(values, module, map[string]map[int]int{module: published})
}
// PreparationArgument is how the mesh asks a module to prepare its state: one word, to the module's
// own program, whatever that program is (novox/hq ADR 0135).
//
// **One word for every kind of module.** A module built as a Go binary receives it as its argument;
// one built as a bundle receives it through the runtime, whose entry takes the same word. So the
// mesh has one way of asking and a module has one way of answering, and neither learns the other's
// shape.
const PreparationArgument = "prepare"
// preparationTarget is the resource a module's preparation runs before: its own workload.
//
// The first container carrying an artifact this module built, and not itself a step — that is the
// thing that runs the module's code, and therefore the thing whose state must be ready. Empty when
// the module prepares nothing, or when nothing it declares could run its code.
//
// **A module with two own workloads gates the first of them.** Five modules in the catalogue declare
// more than one container of their own, none of them preparing anything today. If one ever does and
// its second workload shares the state, the gate is in front of the first — stated here because the
// alternative is a field asking an author to restate what the mesh can see.
func preparationTarget(m Manifest) string {
if !m.Prepares {
return ""
}
for _, r := range m.Resources {
if fmt.Sprint(r["type"]) != "container" || !ownArtifact(r, m.Module) {
continue
}
if once, _ := r["run-once"].(bool); once {
continue
}
return fmt.Sprint(r["id"])
}
return ""
}
// ownArtifact is whether a resource runs something this module built, in either spelling a manifest
// may be in: naming the artifact, before a build resolved it, or carrying the reference a build
// recorded — this mesh's own store, under this module's name.
func ownArtifact(resource map[string]any, module string) bool {
if named, _ := resource["artifact"].(string); named != "" {
return true
}
image, _ := resource["image"].(string)
return strings.HasPrefix(image, ArtifactStoreScheme+module+"/")
}
// prepared is the module's own resource as the step that prepares its state: the same image, the same
// context, run to completion with the mesh's preparation argument.
//
// Three things are taken away rather than copied, each because the step runs while the version it
// prepares for is still running. A published port cannot be bound twice, and a step that tried would
// fail for a reason that has nothing to do with the state. A fixed address cannot be held twice, for
// the same reason. And a cadence is what a step is the opposite of: a container runs once and gates,
// or on a schedule, or stays up, never two (ADR 0053).
func prepared(from map[string]any) map[string]any {
step := map[string]any{}
for k, v := range from {
step[k] = v
}
step["id"] = fmt.Sprint(from["id"]) + ".prepare"
step["name"] = fmt.Sprint(from["name"]) + "-prepare"
step["run-once"] = true
step["args"] = []any{PreparationArgument}
delete(step, "ports")
delete(step, "ip")
delete(step, "schedule")
delete(step, "reload-on")
return step
}
+24
View File
@@ -225,6 +225,19 @@ type Manifest struct {
// subscription to the queue it writes to (design 29 §2).
Uses []string `json:"uses,omitempty"`
// Prepares says this module has state that must be brought to the shape this version needs
// before this version runs, and that the module's own code does it (novox/hq ADR 0135).
//
// **A word, not an arrangement.** The mesh runs the module's own program in its preparation
// mode, in the module's own context — every binding, credential and setting its code receives,
// because it *is* its code. Nothing here names a container, a command, a mount or a variable:
// the module already said all of that once, and a second copy is a second thing to drift.
//
// **Declared, never inferred.** The control plane cannot read what is inside an artifact, so a
// module that ships a migration and does not say this breaks on its first upgrade. That is
// stated in the record rather than guarded here, because nothing mechanical can guard it.
Prepares bool `json:"prepares,omitempty"`
// Tools are the tools this module answers — request and reply, awaited.
//
// **New, and not `serves`**, which this manifest already uses for the facts a consumer needs
@@ -1276,6 +1289,17 @@ func ParseManifest(raw []byte) (Manifest, error) {
"program that reads what the mesh delivered and reconciles",
m.Module, r["id"]))
}
// **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
// resource that runs that program — and a module declaring none has asked for something the mesh
// cannot compose. Said here, where the manifest is read, rather than by a declaration that
// quietly prepares nothing.
if m.Prepares && preparationTarget(m) == "" {
problems = append(problems, fmt.Sprintf(
"%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 "+
"run it in", m.Module))
}
// **A run-once container is a step the host runs to completion** (novox/hq ADR 0052). It is a
// boolean modifier on the container shape — the host runs the container, requires it to exit 0,
// and starts whatever the declaration places after it only once it has. A value that is not a
+135
View File
@@ -0,0 +1,135 @@
package catalogue
import (
"encoding/json"
"fmt"
"testing"
)
// A module version prepares its state before it runs (novox/hq ADR 0135).
//
// What the mesh derives is the module's own resource, run once with one word, placed immediately in
// front of the thing it prepares for. What matters in these tests is that the derivation is a copy
// rather than a second description: the failure it replaces was a hand-written step repeating six
// fields of the resource it preceded, each free to drift from it.
func aPreparingModule() Manifest {
return Manifest{
Module: "gitea",
Prepares: true,
Resources: []map[string]any{
{"id": "state", "type": "directory", "path": "/var/lib/gitea", "mode": "0700"},
{"id": "server", "type": "container", "name": "mesh-gitea-server",
"image": "gitea/gitea@" + digest, "ports": []any{"3000:3000"}},
{"id": "runtime", "type": "container", "name": "mesh-gitea",
"image": ArtifactStoreScheme + "gitea/runtime@" + digest, "network": "host",
"env": map[string]any{"MESH_GITEA_STATE_DIR": "/run/state"},
"volumes": []any{"/var/lib/gitea:/run/state:ro"},
"ports": []any{"9000:9000"}},
},
}
}
func declaredFor(t *testing.T, m Manifest) []map[string]any {
t.Helper()
// The manifest as the mesh holds it: a build resolved the module's own artifact into the
// reference it recorded, which is also how the composition knows whose code a resource runs.
out, err := Resolution{Node: "anchor", Modules: []Manifest{m}}.Declaration(
Rendering{ArtifactStore: "anchor.internal:5100"})
if err != nil {
t.Fatal(err)
}
return out
}
func idsOf(resources []map[string]any) []string {
var ids []string
for _, r := range resources {
ids = append(ids, fmt.Sprint(r["id"]))
}
return ids
}
// The step runs the module's own code, and comes immediately before it — not before the upstream
// server the module packages, which may be the very thing the state lives in.
func TestThePreparationRunsTheModulesOwnCodeAndComesRightBeforeIt(t *testing.T) {
out := declaredFor(t, aPreparingModule())
ids := idsOf(out)
at := -1
for i, id := range ids {
if id == "gitea.runtime.prepare" {
at = i
}
}
if at < 0 {
t.Fatalf("nothing prepares this module's state: %v", ids)
}
if ids[at+1] != "gitea.runtime" {
t.Fatalf("the preparation is not immediately before the module's own code: %v", ids)
}
for _, id := range ids[:at] {
if id == "gitea.runtime" {
t.Fatalf("the module's own code runs before its state is prepared: %v", ids)
}
}
}
// It is given exactly what the module's own code is given. Asserted field by field against the
// resource it was derived from, because writing it twice is the fault this replaces.
func TestThePreparationIsGivenWhatTheModuleIsGiven(t *testing.T) {
out := declaredFor(t, aPreparingModule())
declared := byID(out)
step, workload := declared["gitea.runtime.prepare"], declared["gitea.runtime"]
if step == nil || workload == nil {
t.Fatalf("expected both, got %v", idsOf(out))
}
for _, field := range []string{"image", "network", "env", "volumes", "type"} {
if fmt.Sprint(step[field]) != fmt.Sprint(workload[field]) {
t.Errorf("the preparation's %s is %v and the module's is %v", field, step[field], workload[field])
}
}
if once, _ := step["run-once"].(bool); !once {
t.Error("the preparation is not a step, so nothing waits for it and nothing is gated by it")
}
if fmt.Sprint(step["args"]) != fmt.Sprint([]any{PreparationArgument}) {
t.Errorf("the preparation is asked for as %v", step["args"])
}
if fmt.Sprint(step["name"]) == fmt.Sprint(workload["name"]) {
t.Error("the preparation and the workload have one name, so one removes the other")
}
// A published port cannot be bound twice, and the version being replaced is still running.
if _, published := step["ports"]; published {
t.Errorf("the preparation publishes a port the running version holds: %v", step["ports"])
}
}
// A module that says nothing about preparing gets nothing, which is most modules.
func TestAModuleThatPreparesNothingGetsNoStep(t *testing.T) {
m := aPreparingModule()
m.Prepares = false
for _, id := range idsOf(declaredFor(t, m)) {
if id == "gitea.runtime.prepare" {
t.Fatal("a module that prepares nothing was given a preparation")
}
}
}
// A module whose own code the mesh cannot find has nothing to ask, and saying so where the manifest
// is read beats a declaration that quietly prepares nothing.
func TestAModuleThatPreparesAndRunsNoneOfItsOwnCodeIsRefused(t *testing.T) {
m := Manifest{
Module: "gitea",
Prepares: true,
Resources: []map[string]any{
// Only the upstream server it packages: nothing here runs gitea's own code.
{"id": "server", "type": "container", "name": "mesh-gitea-server", "image": "gitea/gitea@" + digest},
},
}
raw, err := json.Marshal(m)
if err != nil {
t.Fatal(err)
}
if _, err := ParseManifest(raw); err == nil {
t.Fatal("a module that prepares its state with nothing of its own to run was accepted")
}
}
+37
View File
@@ -154,10 +154,47 @@ func (n *natsInbound) deliver(ctx context.Context, act func(context.Context, Con
m.seq = meta.Sequence.Stream
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)
}
// 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.
//
// By subject, which is the only thing the server enforces: a body claiming to be a report does not
+57
View File
@@ -395,3 +395,60 @@ func TestEverySubjectTheControllerFollowsDecodesToAKind(t *testing.T) {
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 }
+1
View File
@@ -27,6 +27,7 @@
"bus": "/var/lib/mesh/mesh-controller/bus"
},
"secrets-owner": "65534:65534",
"prepares": true,
"resources": [
{
"id": "mesh-state",