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/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)) } 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()) + } +}