Compare commits

..
6 Commits
Author SHA1 Message Date
jochen e11caecdad Let two controllers overlap safely while one hands over to the other (hq issue 213)
The controller's machine moves it from the container to a process by
starting the process first and removing the container once the process
is up (mesh-host's `replaces`). For that moment two controllers share the
store and the bus. Checked what each does:

- the seat's verbs: a queue group per seat, each call answered once. Safe.
- the controller's consumers on CONTROL and EVENTS: push consumers with
  no delivery group, so the second bind is refused with "consumer is
  already bound" and serve exited. The process would restart for ever,
  the host would never see it up, and the container would never go. The
  second controller now stands by and binds when the first lets go
  (tested on a real bus; fails without the change).
- plans: read, changed and saved whole by the 30s timer, by build
  outcomes, by a merge and by `plans stop`. Two timers would each ask a
  tier the other had just asked. Working the plans now takes a
  session-level advisory lock on the inventory: the timer skips while
  another holds it, the other paths wait for it. Build asks happen only
  inside plan work and are covered by the same lock.
2026-10-04 01:11:26 +02:00
jochen 7bb9e55d0b Compose a module's Go service as a process the host runs (hq issue 213)
The controller is to be declared as a Go bundle run by a process instead of
an image (novox/hq issue 213, ADR 0188 §1, §3). The composer could not
express that honestly yet:

- a module declaring tools had every bundle served by the node's runtime,
  so the controller's own binary would have been launched a second time as
  an MCP child; a bundle one of the module's resources runs is now served
  only when it says `loads`
- a module's accounts went after the mesh-computed files, so secrets owned
  by the account a process runs as were refused on the first apply; a
  module's `user` resources now go first
- `prepares` derived its step only from a container; a process is now
  prepared by the same program with `prepare` as a run-once process
- a process may say what it `replaces` (a resource of its module it no
  longer declares), prefixed as the host records it, so the host keeps the
  old one running until the process is (needs mesh-host's `replaces`)

This lands before the controller's manifest uses any of it: the running
controller composes its own declaration, so the code that fills the new
shape must be live first.
2026-10-04 01:11:26 +02:00
mesh-admin 00037608ae Merge pull request 'The forge's tests compose its code as a bundle the node's runtime serves (hq ADR 0198, to-be 38 WP4c)' (#254) from feat/0198-waves-2-3-the-forges-code-is-a-bundle into main 2026-10-03 23:01:17 +00:00
jochen cea59428b1 The forge's tests compose its code as a bundle the node's runtime serves (hq ADR 0198)
gitea's own code moves out of its runtime container (mesh-catalog, to-be 38 WP4c waves 2-3), so the three tests that composed the forge from the catalogue beside this checkout resolve its build as the code bundle, compose it beside the node's runtime, and read the forge's address from the words the runtime hands the module rather than from a sidecar's env.
2026-10-04 00:50:06 +02:00
mesh-admin 5c832f2d19 Merge pull request 'Compose a process's environment as a container's' (#250) from feat/a-process-env-is-composed-like-a-containers into main 2026-10-03 22:29:11 +00:00
jochen cdebb7d1a5 Compose a process's environment as a container's
A module's own code moving out of its container (novox/hq to-be 38 WP4c)
becomes a process on the machine, and still has to be told what its
container was: the port this machine gave the module and where the
foundation's seats are. ${port:…} and ${seat:…} were filled only in a
file's content and a container's env, so in a process's env they reached
the machine as literals, and the modules that moved first (mesh-catalog
#245) wrote their run-once steps a 0600 env file instead. A process's env
now takes the same resolution and the same refusals; ${dir:…} and
${access:…} already did, and a bundle's env (ADR 0192) already resolves
${dir:…} and ${port:…}.
2026-10-04 00:27:42 +02:00
12 changed files with 430 additions and 41 deletions
@@ -0,0 +1,47 @@
package main
import (
"testing"
"time"
"github.com/novox/mesh-controller/internal/inventory"
)
// novox/hq issue 213: for the moment a machine hands its controller over, the container and the
// process both run the plan timer on one store. Only the one holding the plans moves them; the other
// leaves them alone, and moves them once they are let go.
func TestAControllerLeavesThePlansToTheOneHoldingThem(t *testing.T) {
open := aMesh(t)
ctx := t.Context()
now := time.Now().UTC()
// Every tier done: the next step is the plan's last, and needs nothing but the store.
plan := inventory.Plan{ID: "plan-213", Repository: "r", Commit: "abc", Created: now, Updated: now,
State: inventory.PlanRolling, Tier: 1, Tiers: [][]string{{"app"}},
Modules: map[string]*inventory.PlanModule{"app": {State: "built"}}}
if err := open.inventory.SavePlan(ctx, plan); err != nil {
t.Fatal(err)
}
// The other controller: its own connections to the same store, holding the plans.
other, err := inventory.Open(ctx)
if err != nil {
t.Fatal(err)
}
t.Cleanup(other.Close)
release, err := other.HoldPlans(ctx, false)
if err != nil {
t.Fatal(err)
}
t.Cleanup(release) // before the close above: a pool waits for a connection still held
advancePlans(ctx, open)
if p, err := open.inventory.PlanByID(ctx, "plan-213"); err != nil || !p.Open() {
t.Fatalf("a controller moved a plan another held: %+v %v", p, err)
}
release()
advancePlans(ctx, open)
if p, err := open.inventory.PlanByID(ctx, "plan-213"); err != nil || p.State != inventory.PlanDone {
t.Fatalf("the plan did not move once it was let go: %+v %v", p, err)
}
}
+41 -1
View File
@@ -2,6 +2,7 @@ package main
import ( import (
"context" "context"
"errors"
"flag" "flag"
"fmt" "fmt"
"sort" "sort"
@@ -296,6 +297,15 @@ func askTier(ctx context.Context, inv *inventory.Inventory, p *inventory.Plan) e
// build's request time is not known, and such an outcome is taken as before. // build's request time is not known, and such an outcome is taken as before.
func planBuilt(ctx context.Context, open *stores, module, commit, failed string, asked time.Time) { func planBuilt(ctx context.Context, open *stores, module, commit, failed string, asked time.Time) {
inv := open.inventory inv := open.inventory
// One controller works the plans at a time (novox/hq issue 213); an outcome waits its turn rather
// than write over what the holder is about to save. Not taken, it is still in the build records,
// which the holder settles the plan from (issue 214).
release, err := inv.HoldPlans(ctx, true)
if err != nil {
fmt.Printf("plans: %s's outcome is left to the build records: %v\n", module, err)
return
}
defer release()
plans, err := inv.OpenPlans(ctx) plans, err := inv.OpenPlans(ctx)
if err != nil { if err != nil {
fmt.Printf("plans: cannot read them: %v\n", err) fmt.Printf("plans: cannot read them: %v\n", err)
@@ -342,13 +352,30 @@ func planBuilt(ctx context.Context, open *stores, module, commit, failed string,
fmt.Printf("%s: %s; the tiers after it are not asked\n", p.ID, p.Note) fmt.Printf("%s: %s; the tiers after it are not asked\n", p.ID, p.Note)
} }
} }
advancePlans(ctx, open) advanceHeld(ctx, open)
} }
// advancePlans moves every open plan as far as the facts allow: a tier whose modules are all built // advancePlans moves every open plan as far as the facts allow: a tier whose modules are all built
// and whose gates are applied gives way to the next; the last tier done is the plan done. Called // and whose gates are applied gives way to the next; the last tier done is the plan done. Called
// after every outcome and on a timer, so a plan waiting on a machine's report moves when it comes. // after every outcome and on a timer, so a plan waiting on a machine's report moves when it comes.
//
// **One controller at a time** (novox/hq issue 213). A plan is read, changed and saved whole; two
// controllers — the old and the new while a machine hands its controller over — would each ask a
// tier the other had just asked. Taken without waiting: whoever holds the plans is moving them.
func advancePlans(ctx context.Context, open *stores) { func advancePlans(ctx context.Context, open *stores) {
release, err := open.inventory.HoldPlans(ctx, false)
if err != nil {
if !errors.Is(err, inventory.ErrPlansBusy) {
fmt.Printf("plans: cannot hold them: %v\n", err)
}
return
}
defer release()
advanceHeld(ctx, open)
}
// advanceHeld is advancePlans for a caller already holding the plans.
func advanceHeld(ctx context.Context, open *stores) {
inv := open.inventory inv := open.inventory
plans, err := inv.OpenPlans(ctx) plans, err := inv.OpenPlans(ctx)
if err != nil { if err != nil {
@@ -676,6 +703,19 @@ func plansCommand(ctx context.Context, args []string) error {
} }
p.State = inventory.PlanFailed p.State = inventory.PlanFailed
p.Note = "stopped by hand at tier " + fmt.Sprint(p.Tier) p.Note = "stopped by hand at tier " + fmt.Sprint(p.Tier)
release, err := inv.HoldPlans(ctx, true)
if err != nil {
return err
}
defer release()
if p, err = inv.PlanByID(ctx, positionals[1]); err != nil {
return err
}
if !p.Open() {
return fmt.Errorf("%s is already %s", p.ID, p.State)
}
p.State = inventory.PlanFailed
p.Note = "stopped by hand at tier " + fmt.Sprint(p.Tier)
if err := inv.SavePlan(ctx, p); err != nil { if err := inv.SavePlan(ctx, p); err != nil {
return err return err
} }
+7
View File
@@ -315,6 +315,13 @@ func (f following) SourceMoved(ctx context.Context, m link.SourceMoved) error {
for _, e := range moved { for _, e := range moved {
movedNames = append(movedNames, e.Manifest.Module) movedNames = append(movedNames, e.Manifest.Module)
} }
// Written and its first tier asked as one act on the plans (novox/hq issue 213): a timer on
// another controller reading it between the two would ask the tier again.
release, err := inv.HoldPlans(ctx, true)
if err != nil {
return notNow(err)
}
defer release()
plan := planOfMerge(m, movedNames, edges) plan := planOfMerge(m, movedNames, edges)
if hasCycle(plan.Tiers, edges) { if hasCycle(plan.Tiers, edges) {
fmt.Printf(" the last tier depends on itself: %s — built together, in no order\n", fmt.Printf(" the last tier depends on itself: %s — built together, in no order\n",
@@ -146,18 +146,6 @@ func TestWhatAProcessReplacesIsNamedAsTheHostRecordedIt(t *testing.T) {
} }
} }
// Where the machine put the store is told to a process as to a container.
func TestAProcessIsToldWhereTheSeatsAre(t *testing.T) {
out := composeTheService(t, Rendering{Seats: map[string]map[int]int{"mesh-store": {5432: 6852}}})
env, _ := out[indexOf(out, "svc.service")]["env"].(map[string]any)
if env["SVC_STORE_PORT"] != "6852" {
t.Fatalf("the process is told the store is on %v; the node put it on 6852", env["SVC_STORE_PORT"])
}
if !strings.HasPrefix(env["SVC_STORE_FILE"].(string), "/") {
t.Errorf("the secret's path was not placed: %v", env["SVC_STORE_FILE"])
}
}
func TestWhatReplacesMayNameIsRefusedNearItsAuthor(t *testing.T) { func TestWhatReplacesMayNameIsRefusedNearItsAuthor(t *testing.T) {
for what, resource := range map[string]string{ for what, resource := range map[string]string{
"a container": `{"id":"c","type":"container","name":"c","image":"x@` + aServiceDigest + `","replaces":["old"]}`, "a container": `{"id":"c","type":"container","name":"c","image":"x@` + aServiceDigest + `","replaces":["old"]}`,
+23 -17
View File
@@ -1,6 +1,7 @@
package catalogue package catalogue
import ( import (
"encoding/json"
"fmt" "fmt"
"os" "os"
"reflect" "reflect"
@@ -170,7 +171,7 @@ func TestTheForgeHoldsTheNpmAndGitSeats(t *testing.T) {
// **And the forge's own address follows it**, composed from the manifest in the catalogue beside // **And the forge's own address follows it**, composed from the manifest in the catalogue beside
// this checkout (novox/hq 04-ISSUES/088). // this checkout (novox/hq 04-ISSUES/088).
// //
// The forge is reached a third way that neither test above covers: by its own sidecar, over the // The forge is reached a third way that neither test above covers: by its own code, over the
// machine's loopback, told where to go in its environment. The `2999:3000` mapping that lets the // machine's loopback, told where to go in its environment. The `2999:3000` mapping that lets the
// forge go on binding 3000 does nothing for a caller dialling the machine — so a literal there is // forge go on binding 3000 does nothing for a caller dialling the machine — so a literal there is
// wrong on every node whose assignment differs, and wrong for a second reason on a node given the // wrong on every node whose assignment differs, and wrong for a second reason on a node given the
@@ -178,13 +179,13 @@ func TestTheForgeHoldsTheNpmAndGitSeats(t *testing.T) {
// in an `env` at all is a declaration, not a manifest. // in an `env` at all is a declaration, not a manifest.
func TestTheForgesOwnAddressFollowsThePortTheNodeGaveIt(t *testing.T) { func TestTheForgesOwnAddressFollowsThePortTheNodeGaveIt(t *testing.T) {
forge, err := catalogueManifest(t, "gitea").Resolve([]Built{{ forge, err := catalogueManifest(t, "gitea").Resolve([]Built{{
Name: "runtime", Kind: ArtifactImage, Name: "code", Kind: ArtifactBundle,
Reference: "registry.example/gitea-runtime@sha256:" + strings.Repeat("a", 64), Reference: ArtifactStoreScheme + "gitea/code/blobs/" + bundleDigest, Digest: bundleDigest,
}}) }})
if err != nil { if err != nil {
t.Fatalf("the forge's manifest does not resolve against its own build: %v", err) t.Fatalf("the forge's manifest does not resolve against its own build: %v", err)
} }
r := Resolution{Node: "anchor", Modules: []Manifest{forge}, Needs: []Needed{ r := Resolution{Node: "anchor", Modules: []Manifest{forge, theRuntime(t)}, Needs: []Needed{
{Name: "postgres-database", For: "gitea", From: "anchor", At: "127.0.0.1", {Name: "postgres-database", For: "gitea", From: "anchor", At: "127.0.0.1",
Serves: map[string]any{"port": float64(5432)}, Sealed: "sealed-db"}, Serves: map[string]any{"port": float64(5432)}, Sealed: "sealed-db"},
{Name: "route", For: "gitea", From: "anchor"}, {Name: "route", For: "gitea", From: "anchor"},
@@ -194,8 +195,8 @@ func TestTheForgesOwnAddressFollowsThePortTheNodeGaveIt(t *testing.T) {
// The number this node was given for the forge — the one the machine it is about to run on // The number this node was given for the forge — the one the machine it is about to run on
// already publishes. // already publishes.
out, err := r.Declaration(Rendering{ out, err := r.Declaration(Rendering{ArtifactStore: "anchor.internal:5101",
Needed: map[string]map[string]string{"gitea": {"broker": "sealed-broker"}}, Needed: map[string]map[string]string{RuntimeModule: {"broker": "sealed-broker"}},
Given: map[string]map[int]int{"gitea": {3000: 2999}}, Given: map[string]map[int]int{"gitea": {3000: 2999}},
}) })
if err != nil { if err != nil {
@@ -210,14 +211,19 @@ func TestTheForgesOwnAddressFollowsThePortTheNodeGaveIt(t *testing.T) {
if published := fmt.Sprint(server["ports"]); !strings.Contains(published, "2999:3000") { if published := fmt.Sprint(server["ports"]); !strings.Contains(published, "2999:3000") {
t.Fatalf("the forge is not published on the port this node gave it: %v", server["ports"]) t.Fatalf("the forge is not published on the port this node gave it: %v", server["ports"])
} }
runtime := fileNamed(out, "gitea.runtime") // The forge's own code runs in the node's runtime (novox/hq ADR 0198), given its words there.
runtime := fileNamed(out, RuntimeModule+"."+RuntimeProcessID())
if runtime == nil { if runtime == nil {
t.Fatalf("the forge's sidecar is not in the declaration: %v", out) t.Fatalf("the node's runtime is not in the declaration: %v", ids(out))
} }
env, _ := runtime["env"].(map[string]any) env, _ := runtime["env"].(map[string]string)
if env["MESH_GITEA_URL"] != "http://127.0.0.1:2999" { var given map[string]map[string]string
t.Fatalf("the forge's sidecar dials %v while the machine publishes the forge on 2999 — "+ if err := json.Unmarshal([]byte(env[RuntimeToolEnv]), &given); err != nil {
"whatever reads it dials a dead port", env["MESH_GITEA_URL"]) t.Fatalf("the runtime's %s is not JSON: %q", RuntimeToolEnv, env[RuntimeToolEnv])
}
if given["gitea"]["MESH_GITEA_URL"] != "http://127.0.0.1:2999" {
t.Fatalf("the forge's code dials %v while the machine publishes the forge on 2999 — "+
"whatever reads it dials a dead port", given["gitea"]["MESH_GITEA_URL"])
} }
} }
@@ -229,13 +235,13 @@ func declaredGiteaSsh(t *testing.T, given map[int]int) map[string]any {
t.Helper() t.Helper()
forge := catalogueManifest(t, "gitea") forge := catalogueManifest(t, "gitea")
resolved, err := forge.Resolve([]Built{{ resolved, err := forge.Resolve([]Built{{
Name: "runtime", Kind: ArtifactImage, Name: "code", Kind: ArtifactBundle,
Reference: "registry.example/gitea-runtime@sha256:" + strings.Repeat("a", 64), Reference: ArtifactStoreScheme + "gitea/code/blobs/" + bundleDigest, Digest: bundleDigest,
}}) }})
if err != nil { if err != nil {
t.Fatalf("the forge's manifest does not resolve against its own build: %v", err) t.Fatalf("the forge's manifest does not resolve against its own build: %v", err)
} }
r := Resolution{Node: "anchor", Modules: []Manifest{resolved}, Needs: []Needed{ r := Resolution{Node: "anchor", Modules: []Manifest{resolved, theRuntime(t)}, Needs: []Needed{
{Name: "postgres-database", For: "gitea", From: "anchor", At: "127.0.0.1", {Name: "postgres-database", For: "gitea", From: "anchor", At: "127.0.0.1",
Serves: map[string]any{"port": float64(5432)}, Sealed: "sealed-db"}, Serves: map[string]any{"port": float64(5432)}, Sealed: "sealed-db"},
{Name: "route", For: "gitea", From: "anchor"}, {Name: "route", For: "gitea", From: "anchor"},
@@ -246,8 +252,8 @@ func declaredGiteaSsh(t *testing.T, given map[int]int) map[string]any {
for k, v := range given { for k, v := range given {
givenPorts[k] = v givenPorts[k] = v
} }
out, err := r.Declaration(Rendering{ out, err := r.Declaration(Rendering{ArtifactStore: "anchor.internal:5101",
Needed: map[string]map[string]string{"gitea": {"broker": "sealed-broker"}}, Needed: map[string]map[string]string{RuntimeModule: {"broker": "sealed-broker"}},
Ports: map[string]map[int]int{"gitea": givenPorts}, Ports: map[string]map[int]int{"gitea": givenPorts},
Given: map[string]map[int]int{"gitea": given}, Given: map[string]map[int]int{"gitea": given},
}) })
+10 -5
View File
@@ -31,7 +31,7 @@ import (
// //
// So a module asks. `${port:8080}` is "the machine-side port you gave me for the 8080 I said I // So a module asks. `${port:8080}` is "the machine-side port you gave me for the 8080 I said I
// listen on", and the module writes that where it would otherwise have written a literal — in a // listen on", and the module writes that where it would otherwise have written a literal — in a
// file's content, or in a value of a container's `env`. // file's content, or in a value of a container's or a process's `env`.
// //
// **The environment is filled by the control plane, exactly as a bound value is.** A port is not // **The environment is filled by the control plane, exactly as a bound value is.** A port is not
// secret — the mesh holds it in the clear — so there is nothing for the host to be the only // secret — the mesh holds it in the clear — so there is nothing for the host to be the only
@@ -64,7 +64,12 @@ func portsUsed(content string) []int {
} }
// portInto replaces a resource's ${port:…} placeholders with what this machine assigned — in a // portInto replaces a resource's ${port:…} placeholders with what this machine assigned — in a
// file's content, and in a value of a container's environment. // file's content, and in a value of a container's or a process's environment.
//
// **A process's environment is a container's** (novox/hq to-be 38 WP4c). A module's code moving out
// of its container becomes a process on the machine and still has to be told what the container
// was told; filled for one kind and not the other, the literal reached the process and was read as
// a port, and the modules that moved first wrote their run-once steps a 0600 env file instead.
// //
// A port the module did not say it listens on is refused, for the same reason a binding's unknown // A port the module did not say it listens on is refused, for the same reason a binding's unknown
// key is: the module is asking about something it never declared, and the answer would be a guess. // key is: the module is asking about something it never declared, and the answer would be a guess.
@@ -84,7 +89,7 @@ func portInto(resource map[string]any, module string, listens []Listening, with
} }
resource["content"] = filled resource["content"] = filled
case "container": case "container", "process":
env, ok := resource["env"].(map[string]any) env, ok := resource["env"].(map[string]any)
if !ok { if !ok {
return nil return nil
@@ -106,8 +111,8 @@ func portInto(resource map[string]any, module string, listens []Listening, with
continue continue
} }
value, err := portsFilledInto(written, value, err := portsFilledInto(written,
fmt.Sprintf("%s's container %s sets %s to something that", fmt.Sprintf("%s's %s %s sets %s to something that",
module, resource["name"], key), module, listens, with) module, resource["type"], resource["name"], key), module, listens, with)
if err != nil { if err != nil {
return err return err
} }
+80
View File
@@ -0,0 +1,80 @@
package catalogue
import (
"strings"
"testing"
)
// **A process's environment is composed as a container's is** (novox/hq to-be 38 WP4c).
//
// A module's code moving out of its container becomes a process on the machine, and what its
// container's environment asked for — the port this machine gave the module, the place it put the
// module's directory — it still has to be told. Filled for a container and not for a process, the
// literal `${port:8080}` reached the process as its environment and was read as a port; the modules
// that moved first wrote their run-once steps an env file instead.
func processModule(env map[string]any) Manifest {
return Manifest{
Module: "showcase",
Listens: []Listening{{Port: 8080, From: FromMesh}},
Resources: []map[string]any{
{"id": "data", "type": "directory", "mode": "0700"},
{"id": "setup", "type": "process", "name": "showcase-setup", "run-once": true,
"run": []any{"/usr/bin/showcase", "setup"}, "env": env},
},
}
}
func TestAProcessIsToldItsPortAndItsPlaceInItsEnvironment(t *testing.T) {
env := map[string]any{
"SHOWCASE_URL": "http://127.0.0.1:${port:8080}",
"SHOWCASE_DATA": "${dir:data}/objects",
"SHOWCASE_DB": "127.0.0.1:${seat:mesh-store:5432}",
"GREETING": "hello",
}
out, err := Resolution{Node: "anchor", Modules: []Manifest{processModule(env)}}.Declaration(Rendering{
Ports: map[string]map[int]int{"showcase": {8080: 21000}},
Seats: map[string]map[int]int{"mesh-store": {5432: 6852}},
})
if err != nil {
t.Fatalf("a process asking for its port and its place does not compose: %v", err)
}
setup := fileNamed(out, "showcase.setup")
if setup == nil {
t.Fatalf("the process is not in the declaration: %v", out)
}
got, _ := setup["env"].(map[string]any)
for key, want := range map[string]string{
"SHOWCASE_URL": "http://127.0.0.1:21000",
"SHOWCASE_DATA": "/var/lib/showcase/data/objects",
"SHOWCASE_DB": "127.0.0.1:6852",
"GREETING": "hello",
} {
if got[key] != want {
t.Errorf("the process is told %s=%v, want %q", key, got[key], want)
}
}
if env["SHOWCASE_URL"] != "http://127.0.0.1:${port:8080}" {
t.Fatalf("composing for one machine edited the module's own manifest: %v", env)
}
}
// An unknown reference in a process's environment is refused as a container's is, naming the
// process and the variable — left alone, it would reach the machine as a literal.
func TestAProcessAskingAboutAnUndeclaredPortIsRefused(t *testing.T) {
env := map[string]any{"SHOWCASE_URL": "http://127.0.0.1:${port:9999}"}
_, err := Resolution{Node: "anchor", Modules: []Manifest{processModule(env)}}.Declaration(Rendering{})
if err == nil {
t.Fatal("a process was told a port its module never said it listens on")
}
for _, said := range []string{"showcase-setup", "SHOWCASE_URL", "${port:9999}", "8080"} {
if !strings.Contains(err.Error(), said) {
t.Errorf("the refusal does not say %q: %v", said, err)
}
}
env = map[string]any{"SHOWCASE_DATA": "${dir:date}/objects"}
if _, err := (Resolution{Node: "anchor", Modules: []Manifest{processModule(env)}}).Declaration(Rendering{}); err == nil ||
!strings.Contains(err.Error(), "${dir:date}") {
t.Fatalf("a process naming no directory of its module was not refused: %v", err)
}
}
+2 -4
View File
@@ -43,8 +43,8 @@ import (
var ofSeat = regexp.MustCompile(`\$\{seat:([a-z0-9][a-z0-9-]*):([0-9]+)\}`) var ofSeat = regexp.MustCompile(`\$\{seat:([a-z0-9][a-z0-9-]*):([0-9]+)\}`)
// seatInto replaces a resource's ${seat:…} placeholders with where this machine put each seat's // seatInto replaces a resource's ${seat:…} placeholders with where this machine put each seat's
// holder — in a file's content, and in a value of a container's or a process's environment. The same two places // holder — in a file's content, and in a value of a container's or a process's environment. The
// portInto fills, for the same reason: they are where a process reads a number from. // same places portInto fills, for the same reason: they are where a program reads a number from.
func seatInto(resource map[string]any, module string, with Rendering) error { func seatInto(resource map[string]any, module string, with Rendering) error {
switch fmt.Sprint(resource["type"]) { switch fmt.Sprint(resource["type"]) {
case "file": case "file":
@@ -58,8 +58,6 @@ func seatInto(resource map[string]any, module string, with Rendering) error {
} }
resource["content"] = filled resource["content"] = filled
// A process's environment as a container's (novox/hq issue 213): the controller reads where its
// store and broker are from it, whichever way the host runs it.
case "container", "process": case "container", "process":
env, ok := resource["env"].(map[string]any) env, ok := resource["env"].(map[string]any)
if !ok { if !ok {
+59
View File
@@ -87,3 +87,62 @@ func (i *Inventory) tryHold(ctx context.Context, sorted []string) (func(), strin
} }
return release, "", nil return release, "", nil
} }
// ErrPlansBusy is the plans held by another act — on a machine replacing its controller, the other
// controller — for longer than a caller waits, or at all for one that does not wait.
var ErrPlansBusy = errors.New("another controller is working the plans")
// HoldPlans makes working the plans one act at a time, across every controller on the store
// (novox/hq issue 213). A plan is read, changed and written whole; two controllers doing that at
// once — the old and the new for the moment a machine hands its controller over, or a controller
// and a person's `plans stop` — each act on what the other has not saved yet: a tier asked twice,
// an outcome written over. A session-level advisory lock on one connection, released by the
// returned function and by the session ending, so a controller that dies holding it holds nothing.
//
// wait false gives ErrPlansBusy at once when another holds them — the timer's way: the holder is
// moving the plans already. wait true looks again every HoldPoll for up to HoldWaitFor — an
// outcome's or a merge's way, which must be written.
func (i *Inventory) HoldPlans(ctx context.Context, wait bool) (func(), error) {
deadline := time.Now().Add(HoldWaitFor)
for {
release, took, err := i.tryLock(ctx, "mesh-plans")
if err != nil || took {
return release, err
}
if !wait || time.Now().After(deadline) {
return nil, ErrPlansBusy
}
select {
case <-ctx.Done():
return nil, ctx.Err()
case <-time.After(HoldPoll):
}
}
}
// tryLock takes one named advisory lock on a connection of its own, or gives the connection back.
func (i *Inventory) tryLock(ctx context.Context, key string) (func(), bool, error) {
conn, err := i.store.Pool().Acquire(ctx)
if err != nil {
return nil, false, err
}
var once sync.Once
release := func() {
once.Do(func() {
if _, err := conn.Exec(context.WithoutCancel(ctx), `select pg_advisory_unlock_all()`); err != nil {
_ = conn.Conn().Close(context.WithoutCancel(ctx))
}
conn.Release()
})
}
var took bool
if err := conn.QueryRow(ctx, `select pg_try_advisory_lock(hashtext($1)::bigint)`, key).Scan(&took); err != nil {
release()
return nil, false, err
}
if !took {
release()
return nil, false, nil
}
return release, true, nil
}
+38
View File
@@ -0,0 +1,38 @@
package inventory
import (
"errors"
"testing"
"time"
)
// novox/hq issue 213: while a machine hands its controller over from the container to the process,
// two controllers run on one store for a moment. Working the plans is one act at a time across them.
func TestThePlansAreWorkedByOneControllerAtATime(t *testing.T) {
first := ForTest(t)
// A second controller: its own connections to the same store.
second, err := Open(t.Context())
if err != nil {
t.Fatal(err)
}
t.Cleanup(second.Close)
release, err := first.HoldPlans(t.Context(), false)
if err != nil {
t.Fatalf("the plans could not be held when nobody held them: %v", err)
}
t.Cleanup(release) // a pool closing waits for a connection still held; release is idempotent
if _, err := second.HoldPlans(t.Context(), false); !errors.Is(err, ErrPlansBusy) {
t.Fatalf("a second controller held the plans while the first did: %v", err)
}
// A waiter gets them once they are let go.
was := HoldPoll
HoldPoll = 10 * time.Millisecond
defer func() { HoldPoll = was }()
go func() { time.Sleep(50 * time.Millisecond); release() }()
again, err := second.HoldPlans(t.Context(), true)
if err != nil {
t.Fatalf("a waiting controller never got the plans once they were let go: %v", err)
}
again()
}
+54 -2
View File
@@ -5,6 +5,7 @@ import (
"encoding/json" "encoding/json"
"errors" "errors"
"fmt" "fmt"
"log"
"strings" "strings"
"time" "time"
@@ -78,10 +79,15 @@ func (n *natsInbound) Receive(ctx context.Context, act func(context.Context, Con
// than creating one here: the consumer is an object with a configuration — ack policy, ack // than creating one here: the consumer is an object with a configuration — ack policy, ack
// wait, redelivery — and a client that creates its own would be a second opinion about it. // wait, redelivery — and a client that creates its own would be a second opinion about it.
control := make(chan *nats.Msg, Prefetch) control := make(chan *nats.Msg, Prefetch)
said, err := js.ChanSubscribe("", control, nats.Bind("CONTROL", broker.ControllerName)) said, err := standingBy(ctx, log.Default(), "CONTROL", func() (*nats.Subscription, error) {
return js.ChanSubscribe("", control, nats.Bind("CONTROL", broker.ControllerName))
})
if err != nil { if err != nil {
return fmt.Errorf("subscribing to what nodes say: %w", err) return fmt.Errorf("subscribing to what nodes say: %w", err)
} }
if said == nil {
return nil // stopped while standing by
}
defer func() { _ = said.Unsubscribe() }() defer func() { _ = said.Unsubscribe() }()
// Heartbeats, on core NATS and off any stream (design 25 §3). Their own subscription because // Heartbeats, on core NATS and off any stream (design 25 §3). Their own subscription because
@@ -97,10 +103,15 @@ func (n *natsInbound) Receive(ctx context.Context, act func(context.Context, Con
var events chan *nats.Msg var events chan *nats.Msg
if len(n.follows) > 0 { if len(n.follows) > 0 {
events = make(chan *nats.Msg, Prefetch) events = make(chan *nats.Msg, Prefetch)
followed, err := js.ChanSubscribe("", events, nats.Bind("EVENTS", broker.ControllerName)) followed, err := standingBy(ctx, log.Default(), "EVENTS", func() (*nats.Subscription, error) {
return js.ChanSubscribe("", events, nats.Bind("EVENTS", broker.ControllerName))
})
if err != nil { if err != nil {
return fmt.Errorf("subscribing to what the catalogue says: %w", err) return fmt.Errorf("subscribing to what the catalogue says: %w", err)
} }
if followed == nil {
return nil
}
defer func() { _ = followed.Unsubscribe() }() defer func() { _ = followed.Unsubscribe() }()
} }
@@ -324,3 +335,44 @@ func (m *natsControl) forget() {
type replyAddressed struct { type replyAddressed struct {
ReplyTo string `json:"reply_to,omitempty"` ReplyTo string `json:"reply_to,omitempty"`
} }
// StandbyPoll is how often a controller standing by looks again for its consumers. A variable so a
// test need not wait.
var StandbyPoll = 2 * time.Second
// standingBy binds one of the controller's consumers, waiting while another controller holds it.
//
// **Two controllers, one consumer** (novox/hq issue 213). The controller's consumers are push
// consumers with no delivery group, so the server lets one subscription bind each — on purpose:
// two would each act on every message (issue 146). When a machine hands its controller over from
// the container to the process, the host starts the process first and removes the container only
// once the process is up; the process then finds the consumers bound. Exiting on that would never
// be up, so the container would never go. It stands by instead — the seat's verbs are already
// served from a queue group, and the plans wait on their lock — and binds as soon as the other lets
// go. Nil and no error is ctx ending while it waited.
func standingBy(ctx context.Context, logger interface{ Printf(string, ...any) }, stream string,
bind func() (*nats.Subscription, error)) (*nats.Subscription, error) {
said := false
for {
sub, err := bind()
if err == nil {
if said {
logger.Printf("took the controller's consumer on %s: the controller that held it let go", stream)
}
return sub, nil
}
if !strings.Contains(err.Error(), "already bound") {
return nil, err
}
if !said {
logger.Printf("another controller holds the controller's consumer on %s; standing by "+
"until it lets go", stream)
said = true
}
select {
case <-ctx.Done():
return nil, nil
case <-time.After(StandbyPoll):
}
}
}
+69
View File
@@ -0,0 +1,69 @@
package link
import (
"context"
"encoding/json"
"os"
"testing"
"time"
"github.com/novox/mesh-controller/internal/broker"
)
// novox/hq issue 213: while a machine hands its controller over, the new controller (the process)
// is started while the old one (the container) still holds the controller's consumers. It must not
// exit — the host would read that as a replacement that did not come up and never remove the
// container — and must not act on what the old one is handed. It stands by, and takes the consumers
// when the old one lets go.
func TestNatsASecondControllerStandsByAndTakesOverWhenTheFirstLetsGo(t *testing.T) {
js := aBus(t)
was := StandbyPoll
StandbyPoll = 50 * time.Millisecond
defer func() { StandbyPoll = was }()
old := &counted{}
_, stopOld := servingOn(t, js, old)
eventually(t, "the first controller binding its consumer", func() bool {
info, err := js.Context().ConsumerInfo("CONTROL", broker.ControllerName)
return err == nil && info.PushBound
})
// The new one, on a connection of its own as the process would have.
second, err := broker.Dial(os.Getenv("MESH_TEST_NATS"))
if err != nil {
t.Fatal(err)
}
t.Cleanup(second.Close)
fresh := &counted{}
s := &Server{inbound: Nats(second), bus: OverNATS{Conn: second.Conn(), JS: second.Context()},
listener: fresh, log: quiet()}
ctx, stopNew := context.WithCancel(context.Background())
defer stopNew()
ended := make(chan error, 1)
go func() { ended <- s.Serve(ctx) }()
select {
case err := <-ended:
t.Fatalf("the second controller stopped instead of standing by: %v", err)
case <-time.After(500 * time.Millisecond):
}
report := func(declared string) {
body, _ := json.Marshal(Report{Node: "anchor", Declared: declared, Applied: []string{"store"}})
if _, err := js.Context().Publish(ReportSubject("anchor"), body); err != nil {
t.Fatal(err)
}
}
report("d1")
eventually(t, "the holding controller hearing the report", func() bool { return old.count() == 1 })
if fresh.count() != 0 {
t.Fatal("the controller standing by acted on a report the holder was handed")
}
stopOld()
report("d2")
eventually(t, "the second controller taking over once the first let go", func() bool { return fresh.count() == 1 })
if old.count() != 1 {
t.Errorf("the first controller heard %d reports", old.count())
}
}