From 7bb9e55d0bc2a445825ae5208ef33d973bce46e2 Mon Sep 17 00:00:00 2001 From: jochen Date: Sun, 4 Oct 2026 01:02:25 +0200 Subject: [PATCH 1/2] Compose a module's Go service as a process the host runs (hq issue 213) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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. --- internal/catalogue/a_service_process_test.go | 166 +++++++++++++++++++ internal/catalogue/build.go | 27 ++- internal/catalogue/declaration.go | 56 ++++++- internal/catalogue/manifest.go | 49 +++++- 4 files changed, 286 insertions(+), 12 deletions(-) create mode 100644 internal/catalogue/a_service_process_test.go diff --git a/internal/catalogue/a_service_process_test.go b/internal/catalogue/a_service_process_test.go new file mode 100644 index 0000000..1e25276 --- /dev/null +++ b/internal/catalogue/a_service_process_test.go @@ -0,0 +1,166 @@ +package catalogue + +import ( + "encoding/json" + "strings" + "testing" +) + +// novox/hq issue 213 (ADR 0188 §1, §3): a module's own Go service is a bundle the host runs as a +// process, not an image. Each test holds one thing that had to change in the composer for the +// controller to be declared that way. + +const aServiceDigest = "sha256:dddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddd" + +// aServiceModule is the controller's shape in miniature: it answers tools of its own, its code is a +// Go bundle a process runs as an account it declares, its secrets belong to that account, it +// prepares its state, and its process replaces the container it used to run as. +func aServiceModule(t *testing.T) Manifest { + t.Helper() + raw := `{ + "module": "svc", "version": "1", "tools": ["status"], "prepares": true, + "own-secrets": {"store": "${dir:state}/store"}, + "secrets-owner": "svc", + "resources": [ + {"id": "state", "type": "directory", "mode": "0700", "place": "mesh", "owner": "svc"}, + {"id": "service", "type": "process", "name": "svc", "artifact": "code", + "run": ["./svc", "serve"], "user": "svc", "replaces": ["server"], + "env": {"SVC_STORE_FILE": "${dir:state}/store", "SVC_STORE_PORT": "${seat:mesh-store:5432}"}}, + {"id": "account", "type": "user", "name": "svc", "shell": "/usr/bin/nologin", "home": "/var/lib/svc"} + ], + "build": {"artifacts": [{"name": "code", "kind": "bundle", "language": "go", "system": "arch", + "from": "cmd/svc", "binary": "svc"}]} + }` + m, err := ParseManifest([]byte(raw)) + if err != nil { + t.Fatalf("the service's manifest is refused: %v", err) + } + resolved, err := m.Resolve([]Built{{Name: "code", Kind: ArtifactBundle, + Reference: ArtifactStoreScheme + "svc/code@" + aServiceDigest, Digest: aServiceDigest}}) + if err != nil { + t.Fatal(err) + } + return resolved +} + +func composeTheService(t *testing.T, with Rendering) []map[string]any { + t.Helper() + with.Needed = map[string]map[string]string{"svc": {"store": "sealed-store"}} + with.ArtifactStore = "anchor.internal:5100" + out, err := Resolution{Node: "anchor", Modules: []Manifest{aServiceModule(t)}}.Declaration(with) + if err != nil { + t.Fatalf("the service does not compose: %v", err) + } + return out +} + +func indexOf(out []map[string]any, id string) int { + for i, r := range out { + if r["id"] == id { + return i + } + } + return -1 +} + +// A module that declares tools has every bundle served by the node's runtime unless it says +// otherwise — and the controller declares the verbs it answers as tools. Its service bundle is run +// by its own process; launched a second time by the runtime it would be a second controller +// pretending to be an MCP server. +func TestABundleItsOwnProcessRunsIsNotServedByTheRuntime(t *testing.T) { + m := aServiceModule(t) + if len(m.Bundles) != 1 { + t.Fatalf("the service's bundle was not kept: %+v", m.Bundles) + } + if loads := m.Bundles[0].Loads; len(loads) != 0 { + t.Fatalf("the runtime would launch the service's own bundle as tools: %v", loads) + } + // And a bundle no resource runs still is served, as a module declaring tools always had it. + tools := Manifest{Module: "t", Version: "1", Tools: []string{"x"}, + Build: &Build{Artifacts: []Artifact{{Name: "tools", Kind: ArtifactBundle, Language: "go", + System: "arch", From: "cmd/t"}}}} + resolved, err := tools.Resolve([]Built{{Name: "tools", Kind: ArtifactBundle, + Reference: ArtifactStoreScheme + "t/tools@" + aServiceDigest, Digest: aServiceDigest}}) + if err != nil { + t.Fatal(err) + } + if loads := resolved.Bundles[0].Loads; len(loads) != 1 || loads[0] != "t" { + t.Fatalf("a tools bundle nothing runs is no longer served: %v", loads) + } +} + +// The account is created before anything is given to it. Its secrets are mesh-computed and so +// placed before the module's own resources; given to a user the machine did not have yet, they were +// refused on the first apply and the process started without them. +func TestAModulesAccountComesBeforeWhatBelongsToIt(t *testing.T) { + out := composeTheService(t, Rendering{}) + account, secret := indexOf(out, "svc.account"), indexOf(out, "svc."+NeedID("store")) + if account < 0 || secret < 0 { + t.Fatalf("the account or the secret is missing: %v", out) + } + if account > secret { + t.Fatalf("the secret owned by svc is written before svc exists: account at %d, secret at %d", + account, secret) + } + if owner := out[secret]["owner"]; owner != "svc" { + t.Errorf("the secret belongs to %v, not the account its process runs as", owner) + } +} + +// The process is the module's program; its preparation is the same program asked to prepare, as a +// step before it — with the same account and environment, and handing nothing over. +func TestAProcessIsPreparedByItsOwnProgram(t *testing.T) { + out := composeTheService(t, Rendering{}) + step, process := indexOf(out, "svc.service-prepare"), indexOf(out, "svc.service") + if step < 0 || process < 0 || step > process { + t.Fatalf("the preparation is not a step before the process (%d, %d): %v", step, process, out) + } + s := out[step] + if s["type"] != "process" || s["run-once"] != true || s["name"] != "svc-prepare" { + t.Errorf("the preparation is not a run-once process: %v", s) + } + if run, _ := json.Marshal(s["run"]); string(run) != `["./svc","prepare"]` { + t.Errorf("the preparation runs %s", run) + } + if s["user"] != "svc" || s["source"] != out[process]["source"] { + t.Errorf("the preparation does not run the same bundle as the same account: %v", s) + } + if env, _ := s["env"].(map[string]any); env["SVC_STORE_FILE"] == nil { + t.Errorf("the preparation is not given the process's environment: %v", s["env"]) + } + if _, has := s["replaces"]; has { + t.Errorf("the preparation would hand over what the process replaces: %v", s) + } + if _, has := s["args"]; has { + t.Errorf("the preparation carries a container's args: %v", s) + } +} + +// What the process replaces is named as the host recorded it, `.`; unprefixed, the host +// matches nothing and removes the container first, as before. +func TestWhatAProcessReplacesIsNamedAsTheHostRecordedIt(t *testing.T) { + out := composeTheService(t, Rendering{}) + p := out[indexOf(out, "svc.service")] + if got, _ := json.Marshal(p["replaces"]); string(got) != `["svc.server"]` { + t.Fatalf("the process replaces %s", got) + } +} + +func TestWhatReplacesMayNameIsRefusedNearItsAuthor(t *testing.T) { + for what, resource := range map[string]string{ + "a container": `{"id":"c","type":"container","name":"c","image":"x@` + aServiceDigest + `","replaces":["old"]}`, + "a step": `{"id":"p","type":"process","name":"p","run":["./p"],"run-once":true,"replaces":["old"]}`, + "something declared": `{"id":"p","type":"process","name":"p","run":["./p"],"replaces":["p"]}`, + "another module's": `{"id":"p","type":"process","name":"p","run":["./p"],"replaces":["other.old"]}`, + "not a list": `{"id":"p","type":"process","name":"p","run":["./p"],"replaces":"old"}`, + } { + raw := `{"module":"m","version":"1","resources":[` + resource + `]}` + if _, err := ParseManifest([]byte(raw)); err == nil || !strings.Contains(err.Error(), "replace") { + t.Errorf("replaces on %s was accepted: %v", what, err) + } + } + ok := `{"module":"m","version":"1","resources":[{"id":"p","type":"process","name":"p","run":["./p"],"replaces":["old"]}]}` + if _, err := ParseManifest([]byte(ok)); err != nil { + t.Errorf("a process replacing what its module no longer declares was refused: %v", err) + } +} diff --git a/internal/catalogue/build.go b/internal/catalogue/build.go index 4db55a1..7a97b3f 100644 --- a/internal/catalogue/build.go +++ b/internal/catalogue/build.go @@ -77,6 +77,7 @@ func (m Manifest) Resolve(built []Built) (Manifest, error) { // build compare equal. out.Bundles = nil if m.Build != nil { + run := runByAResource(m) for _, a := range m.Build.Artifacts { if a.Kind != ArtifactBundle { continue @@ -84,8 +85,13 @@ func (m Manifest) Resolve(built []Built) (Manifest, error) { made := by[a.Name] // What the runtime loads: what the artifact said, else every entrypoint of a module // that declares tools, else nothing (the field's own rule; see Artifact.Loads). + // + // **Never, unasked, a bundle one of the module's own resources runs** (novox/hq issue 213). + // A process the host runs is the module's service, not its tools: the controller declares + // the verbs it answers as `tools` and serves them itself, and its bundle would otherwise + // have been launched a second time by the node's runtime, as an MCP child it is not. loads := append([]string(nil), a.Loads...) - if a.Loads == nil && len(m.Tools) > 0 { + if a.Loads == nil && len(m.Tools) > 0 && !run[a.Name] { loads = append([]string(nil), a.Entrypoints...) // A bundle compiled to a binary has no entrypoints: the binary is what it is, and what // the runtime starts to serve it (novox/hq ADR 0193). So a Go tools bundle is served @@ -453,6 +459,18 @@ func BinaryOf(a Artifact) string { return a.Name } +// runByAResource is the artifacts one of a module's own resources names — a process that runs it, +// a step, an archive that unpacks it — by name. +func runByAResource(m Manifest) map[string]bool { + named := map[string]bool{} + for _, r := range m.Resources { + if a, ok := r["artifact"].(string); ok && a != "" { + named[a] = true + } + } + return named +} + // undeliveredBundles says which of a module's bundles nothing would ever put on a machine (novox/hq // 04-ISSUES/216). A bundle reaches a machine three ways: the node's runtime serves it (it says // `loads`, or its module declares `tools`), a resource names it (a process, a step, an archive), or @@ -463,12 +481,7 @@ func undeliveredBundles(m Manifest) []string { if m.Build == nil || m.Module == RuntimeModule { return nil } - named := map[string]bool{} - for _, r := range m.Resources { - if a, ok := r["artifact"].(string); ok && a != "" { - named[a] = true - } - } + named := runByAResource(m) var problems []string for _, a := range m.Build.Artifacts { if a.Kind != ArtifactBundle || named[a.Name] || len(a.Loads) > 0 || len(m.Tools) > 0 { diff --git a/internal/catalogue/declaration.go b/internal/catalogue/declaration.go index dbd83f6..d69eb7a 100644 --- a/internal/catalogue/declaration.go +++ b/internal/catalogue/declaration.go @@ -693,7 +693,15 @@ func (r Resolution) compose(with Rendering, owner map[string]string, // Now, and not before: a module whose resources are computed replaces them wholesale, and // merging earlier would throw away the files it still needs. - resources = append(append([]map[string]any{}, first...), resources...) + // + // **Except the module's own accounts, which go before even those** (novox/hq issue 213). What + // the mesh computes may belong to one: a module whose code runs as an account it declares has + // its secrets written owned by that account, and a file given to a user the machine does not + // have yet fails — so on the first apply the secrets were refused, the process started without + // them, and the second apply healed it, which is the fault the paragraph above describes. + // An account depends on nothing the mesh computes. + accounts, rest := accountsFirst(resources) + resources = append(append(accounts, first...), rest...) // No container is given the mesh's names (novox/hq ADR 0148). It used to be: every // container got the whole roster as `--add-host` entries at creation, and a name that @@ -872,6 +880,12 @@ func (r Resolution) compose(with Rendering, owner map[string]string, if renamed := reflectsRenamed(m.Module, resource["reload-on"]); renamed != nil { copied["reload-on"] = renamed } + // And what a process replaces (novox/hq issue 213): a resource of this module's that it + // no longer declares, named as the host recorded it, or the host hands nothing over and + // removes it first. + if renamed := reflectsRenamed(m.Module, resource["replaces"]); renamed != nil { + copied["replaces"] = renamed + } // **What reads one of this module's own secrets is restarted when it changes** (novox/hq // issue 203, issue 206). A credential is re-issued by the mesh, and a container that // mounted the old file keeps the old one open: the build machine ran for an hour on a @@ -2010,12 +2024,18 @@ func preparationTarget(m Manifest) string { return "" } for _, r := range m.Resources { - if fmt.Sprint(r["type"]) != "container" || !ownArtifact(r, m.Module) { + // A container, or a process the host runs from a bundle the module built (novox/hq issue + // 213): the same program in the same context, hosted as a unit rather than a container. + kind := fmt.Sprint(r["type"]) + if (kind != "container" && kind != "process") || !ownArtifact(r, m.Module) { continue } if once, _ := r["run-once"].(bool); once { continue } + if r["schedule"] != nil { + continue + } return fmt.Sprint(r["id"]) } return "" @@ -2029,7 +2049,10 @@ func ownArtifact(resource map[string]any, module string) bool { return true } image, _ := resource["image"].(string) - return strings.HasPrefix(image, ArtifactStoreScheme+module+"/") + // A process or an archive carries what was built as its source (novox/hq issue 213). + source, _ := resource["source"].(string) + return strings.HasPrefix(image, ArtifactStoreScheme+module+"/") || + strings.HasPrefix(source, ArtifactStoreScheme+module+"/") } // prepared is the module's own resource as the step that prepares its state: the same image, the same @@ -2051,7 +2074,19 @@ func prepared(from map[string]any) map[string]any { step["id"] = fmt.Sprint(from["id"]) + "-prepare" step["name"] = fmt.Sprint(from["name"]) + "-prepare" step["run-once"] = true - step["args"] = []any{PreparationArgument} + if fmt.Sprint(from["type"]) == "process" { + // A process says its whole command: the program, then its arguments. The step is the same + // program asked to prepare (novox/hq issue 213). It replaces nothing — what the process + // replaces is handed over to the process, never to the step that runs before it — and a + // step is not restarted, it runs again when what it reads changed, which `restart-on` says. + run := stringsIn(from["run"]) + if len(run) > 0 { + step["run"] = []any{run[0], PreparationArgument} + } + delete(step, "replaces") + } else { + step["args"] = []any{PreparationArgument} + } delete(step, "ports") delete(step, "ip") delete(step, "schedule") @@ -2196,3 +2231,16 @@ func withRestartOn(have any, add []string) []any { } return out } + +// accountsFirst splits a module's resources into its accounts and everything else, each in the order +// written. +func accountsFirst(resources []map[string]any) (accounts, rest []map[string]any) { + for _, r := range resources { + if fmt.Sprint(r["type"]) == "user" { + accounts = append(accounts, r) + continue + } + rest = append(rest, r) + } + return accounts, rest +} diff --git a/internal/catalogue/manifest.go b/internal/catalogue/manifest.go index 4c55564..8528c9f 100644 --- a/internal/catalogue/manifest.go +++ b/internal/catalogue/manifest.go @@ -1515,6 +1515,53 @@ func ParseManifest(raw []byte) (Manifest, error) { "program that reads what the mesh delivered and reconciles", m.Module, r["id"])) } + // **What a process replaces is something the module no longer declares** (novox/hq issue 213). + // The host keeps it running until the process is, then removes it: so it is named by the id the + // module used to give it, it is never a resource the module still declares — that would be + // applied and removed by one declaration — and only a process that stays up has anything to + // hand over to. Said here, near the author, as the host would refuse it far away. + ids := map[string]bool{} + for _, r := range m.Resources { + ids[fmt.Sprint(r["id"])] = true + } + for _, r := range m.Resources { + raw, present := r["replaces"] + if !present { + continue + } + if fmt.Sprint(r["type"]) != "process" { + problems = append(problems, fmt.Sprintf( + "%s: %v says what it replaces, and only a process does", m.Module, r["id"])) + continue + } + if once, _ := r["run-once"].(bool); once || r["schedule"] != nil { + problems = append(problems, fmt.Sprintf( + "%s: %v replaces something and runs once or on a schedule — only a process that stays "+ + "up is there a moment later to hand over to", m.Module, r["id"])) + } + list, ok := raw.([]any) + if !ok { + problems = append(problems, fmt.Sprintf( + "%s: %v replaces %v; replaces is a list of the ids this module no longer declares", + m.Module, r["id"], raw)) + continue + } + for _, item := range list { + id, ok := item.(string) + switch { + case !ok || strings.TrimSpace(id) == "": + problems = append(problems, fmt.Sprintf( + "%s: %v replaces %v, which is not an id", m.Module, r["id"], item)) + case strings.Contains(id, "."): + problems = append(problems, fmt.Sprintf( + "%s: %v replaces %q; a process replaces only a resource of its own module, named "+ + "by its own id", m.Module, r["id"], id)) + case ids[id]: + problems = append(problems, fmt.Sprintf( + "%s: %v replaces %q, which this module still declares", m.Module, r["id"], id)) + } + } + } // **A module that prepares its state must have code the mesh can run** (novox/hq ADR 0135). The // 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 @@ -1522,7 +1569,7 @@ func ParseManifest(raw []byte) (Manifest, error) { // 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 — "+ + "%s says it prepares its state, and declares no container or process 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)) } From e11caecdadeb5de3f33458631ac8d521bfeb2a0e Mon Sep 17 00:00:00 2001 From: jochen Date: Sun, 4 Oct 2026 01:11:26 +0200 Subject: [PATCH 2/2] 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. --- .../plans_one_at_a_time_test.go | 47 +++++++++++++ cmd/mesh-controller/release_plan.go | 42 ++++++++++- cmd/mesh-controller/upgrades.go | 7 ++ internal/inventory/hold.go | 59 ++++++++++++++++ internal/inventory/hold_plans_test.go | 38 ++++++++++ internal/link/receive_nats.go | 56 ++++++++++++++- internal/link/standby_nats_test.go | 69 +++++++++++++++++++ 7 files changed, 315 insertions(+), 3 deletions(-) create mode 100644 cmd/mesh-controller/plans_one_at_a_time_test.go create mode 100644 internal/inventory/hold_plans_test.go create mode 100644 internal/link/standby_nats_test.go diff --git a/cmd/mesh-controller/plans_one_at_a_time_test.go b/cmd/mesh-controller/plans_one_at_a_time_test.go new file mode 100644 index 0000000..d76283b --- /dev/null +++ b/cmd/mesh-controller/plans_one_at_a_time_test.go @@ -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) + } +} diff --git a/cmd/mesh-controller/release_plan.go b/cmd/mesh-controller/release_plan.go index ef8cf64..267d2c7 100644 --- a/cmd/mesh-controller/release_plan.go +++ b/cmd/mesh-controller/release_plan.go @@ -2,6 +2,7 @@ package main import ( "context" + "errors" "flag" "fmt" "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. func planBuilt(ctx context.Context, open *stores, module, commit, failed string, asked time.Time) { 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) if err != nil { 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) } } - advancePlans(ctx, open) + advanceHeld(ctx, open) } // 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 // 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) { + 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 plans, err := inv.OpenPlans(ctx) if err != nil { @@ -676,6 +703,19 @@ func plansCommand(ctx context.Context, args []string) error { } p.State = inventory.PlanFailed 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 { return err } diff --git a/cmd/mesh-controller/upgrades.go b/cmd/mesh-controller/upgrades.go index b9d9cd2..b16542b 100644 --- a/cmd/mesh-controller/upgrades.go +++ b/cmd/mesh-controller/upgrades.go @@ -315,6 +315,13 @@ func (f following) SourceMoved(ctx context.Context, m link.SourceMoved) error { for _, e := range moved { 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) if hasCycle(plan.Tiers, edges) { fmt.Printf(" the last tier depends on itself: %s — built together, in no order\n", diff --git a/internal/inventory/hold.go b/internal/inventory/hold.go index 4ddf382..30f20ac 100644 --- a/internal/inventory/hold.go +++ b/internal/inventory/hold.go @@ -87,3 +87,62 @@ func (i *Inventory) tryHold(ctx context.Context, sorted []string) (func(), strin } 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 +} diff --git a/internal/inventory/hold_plans_test.go b/internal/inventory/hold_plans_test.go new file mode 100644 index 0000000..854451a --- /dev/null +++ b/internal/inventory/hold_plans_test.go @@ -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() +} diff --git a/internal/link/receive_nats.go b/internal/link/receive_nats.go index a1889ee..5472069 100644 --- a/internal/link/receive_nats.go +++ b/internal/link/receive_nats.go @@ -5,6 +5,7 @@ import ( "encoding/json" "errors" "fmt" + "log" "strings" "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 // wait, redelivery — and a client that creates its own would be a second opinion about it. 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 { return fmt.Errorf("subscribing to what nodes say: %w", err) } + if said == nil { + return nil // stopped while standing by + } defer func() { _ = said.Unsubscribe() }() // 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 if len(n.follows) > 0 { 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 { return fmt.Errorf("subscribing to what the catalogue says: %w", err) } + if followed == nil { + return nil + } defer func() { _ = followed.Unsubscribe() }() } @@ -324,3 +335,44 @@ func (m *natsControl) forget() { type replyAddressed struct { 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): + } + } +} diff --git a/internal/link/standby_nats_test.go b/internal/link/standby_nats_test.go new file mode 100644 index 0000000..20e1d2b --- /dev/null +++ b/internal/link/standby_nats_test.go @@ -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()) + } +}