Merge pull request 'Compose a module's Go service as a process the host runs (hq issue 213, 1 of 2)' (#252) from fix/issue-213-the-controller-is-a-process into main
This commit was merged in pull request #252.
This commit is contained in:
@@ -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)
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -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
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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",
|
||||||
|
|||||||
@@ -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, `<module>.<id>`; 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)
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -77,6 +77,7 @@ func (m Manifest) Resolve(built []Built) (Manifest, error) {
|
|||||||
// build compare equal.
|
// build compare equal.
|
||||||
out.Bundles = nil
|
out.Bundles = nil
|
||||||
if m.Build != nil {
|
if m.Build != nil {
|
||||||
|
run := runByAResource(m)
|
||||||
for _, a := range m.Build.Artifacts {
|
for _, a := range m.Build.Artifacts {
|
||||||
if a.Kind != ArtifactBundle {
|
if a.Kind != ArtifactBundle {
|
||||||
continue
|
continue
|
||||||
@@ -84,8 +85,13 @@ func (m Manifest) Resolve(built []Built) (Manifest, error) {
|
|||||||
made := by[a.Name]
|
made := by[a.Name]
|
||||||
// What the runtime loads: what the artifact said, else every entrypoint of a module
|
// 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).
|
// 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...)
|
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...)
|
loads = append([]string(nil), a.Entrypoints...)
|
||||||
// A bundle compiled to a binary has no entrypoints: the binary is what it is, and what
|
// 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
|
// 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
|
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
|
// 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
|
// 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
|
// `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 {
|
if m.Build == nil || m.Module == RuntimeModule {
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
named := map[string]bool{}
|
named := runByAResource(m)
|
||||||
for _, r := range m.Resources {
|
|
||||||
if a, ok := r["artifact"].(string); ok && a != "" {
|
|
||||||
named[a] = true
|
|
||||||
}
|
|
||||||
}
|
|
||||||
var problems []string
|
var problems []string
|
||||||
for _, a := range m.Build.Artifacts {
|
for _, a := range m.Build.Artifacts {
|
||||||
if a.Kind != ArtifactBundle || named[a.Name] || len(a.Loads) > 0 || len(m.Tools) > 0 {
|
if a.Kind != ArtifactBundle || named[a.Name] || len(a.Loads) > 0 || len(m.Tools) > 0 {
|
||||||
|
|||||||
@@ -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
|
// Now, and not before: a module whose resources are computed replaces them wholesale, and
|
||||||
// merging earlier would throw away the files it still needs.
|
// 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
|
// 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
|
// 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 {
|
if renamed := reflectsRenamed(m.Module, resource["reload-on"]); renamed != nil {
|
||||||
copied["reload-on"] = renamed
|
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
|
// **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
|
// 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
|
// 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 ""
|
return ""
|
||||||
}
|
}
|
||||||
for _, r := range m.Resources {
|
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
|
continue
|
||||||
}
|
}
|
||||||
if once, _ := r["run-once"].(bool); once {
|
if once, _ := r["run-once"].(bool); once {
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
|
if r["schedule"] != nil {
|
||||||
|
continue
|
||||||
|
}
|
||||||
return fmt.Sprint(r["id"])
|
return fmt.Sprint(r["id"])
|
||||||
}
|
}
|
||||||
return ""
|
return ""
|
||||||
@@ -2029,7 +2049,10 @@ func ownArtifact(resource map[string]any, module string) bool {
|
|||||||
return true
|
return true
|
||||||
}
|
}
|
||||||
image, _ := resource["image"].(string)
|
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
|
// 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["id"] = fmt.Sprint(from["id"]) + "-prepare"
|
||||||
step["name"] = fmt.Sprint(from["name"]) + "-prepare"
|
step["name"] = fmt.Sprint(from["name"]) + "-prepare"
|
||||||
step["run-once"] = true
|
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, "ports")
|
||||||
delete(step, "ip")
|
delete(step, "ip")
|
||||||
delete(step, "schedule")
|
delete(step, "schedule")
|
||||||
@@ -2196,3 +2231,16 @@ func withRestartOn(have any, add []string) []any {
|
|||||||
}
|
}
|
||||||
return out
|
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
|
||||||
|
}
|
||||||
|
|||||||
@@ -1515,6 +1515,53 @@ func ParseManifest(raw []byte) (Manifest, error) {
|
|||||||
"program that reads what the mesh delivered and reconciles",
|
"program that reads what the mesh delivered and reconciles",
|
||||||
m.Module, r["id"]))
|
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
|
// **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
|
// 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
|
// 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.
|
// quietly prepares nothing.
|
||||||
if m.Prepares && preparationTarget(m) == "" {
|
if m.Prepares && preparationTarget(m) == "" {
|
||||||
problems = append(problems, fmt.Sprintf(
|
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 "+
|
"the preparation is this module's own program, so there has to be one for the mesh to "+
|
||||||
"run it in", m.Module))
|
"run it in", m.Module))
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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
|
||||||
|
}
|
||||||
|
|||||||
@@ -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()
|
||||||
|
}
|
||||||
@@ -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):
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -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())
|
||||||
|
}
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user