Compose a module's Go service as a process the host runs (hq issue 213, 1 of 2) #252

Merged
mesh-admin merged 2 commits from fix/issue-213-the-controller-is-a-process into main 2026-10-03 23:40:38 +00:00
7 changed files with 315 additions and 3 deletions
Showing only changes of commit e11caecdad - Show all commits
@@ -0,0 +1,47 @@
package main
import (
"testing"
"time"
"github.com/novox/mesh-controller/internal/inventory"
)
// novox/hq issue 213: for the moment a machine hands its controller over, the container and the
// process both run the plan timer on one store. Only the one holding the plans moves them; the other
// leaves them alone, and moves them once they are let go.
func TestAControllerLeavesThePlansToTheOneHoldingThem(t *testing.T) {
open := aMesh(t)
ctx := t.Context()
now := time.Now().UTC()
// Every tier done: the next step is the plan's last, and needs nothing but the store.
plan := inventory.Plan{ID: "plan-213", Repository: "r", Commit: "abc", Created: now, Updated: now,
State: inventory.PlanRolling, Tier: 1, Tiers: [][]string{{"app"}},
Modules: map[string]*inventory.PlanModule{"app": {State: "built"}}}
if err := open.inventory.SavePlan(ctx, plan); err != nil {
t.Fatal(err)
}
// The other controller: its own connections to the same store, holding the plans.
other, err := inventory.Open(ctx)
if err != nil {
t.Fatal(err)
}
t.Cleanup(other.Close)
release, err := other.HoldPlans(ctx, false)
if err != nil {
t.Fatal(err)
}
t.Cleanup(release) // before the close above: a pool waits for a connection still held
advancePlans(ctx, open)
if p, err := open.inventory.PlanByID(ctx, "plan-213"); err != nil || !p.Open() {
t.Fatalf("a controller moved a plan another held: %+v %v", p, err)
}
release()
advancePlans(ctx, open)
if p, err := open.inventory.PlanByID(ctx, "plan-213"); err != nil || p.State != inventory.PlanDone {
t.Fatalf("the plan did not move once it was let go: %+v %v", p, err)
}
}
+41 -1
View File
@@ -2,6 +2,7 @@ package main
import ( import (
"context" "context"
"errors"
"flag" "flag"
"fmt" "fmt"
"sort" "sort"
@@ -296,6 +297,15 @@ func askTier(ctx context.Context, inv *inventory.Inventory, p *inventory.Plan) e
// build's request time is not known, and such an outcome is taken as before. // build's request time is not known, and such an outcome is taken as before.
func planBuilt(ctx context.Context, open *stores, module, commit, failed string, asked time.Time) { func planBuilt(ctx context.Context, open *stores, module, commit, failed string, asked time.Time) {
inv := open.inventory inv := open.inventory
// One controller works the plans at a time (novox/hq issue 213); an outcome waits its turn rather
// than write over what the holder is about to save. Not taken, it is still in the build records,
// which the holder settles the plan from (issue 214).
release, err := inv.HoldPlans(ctx, true)
if err != nil {
fmt.Printf("plans: %s's outcome is left to the build records: %v\n", module, err)
return
}
defer release()
plans, err := inv.OpenPlans(ctx) plans, err := inv.OpenPlans(ctx)
if err != nil { if err != nil {
fmt.Printf("plans: cannot read them: %v\n", err) fmt.Printf("plans: cannot read them: %v\n", err)
@@ -342,13 +352,30 @@ func planBuilt(ctx context.Context, open *stores, module, commit, failed string,
fmt.Printf("%s: %s; the tiers after it are not asked\n", p.ID, p.Note) fmt.Printf("%s: %s; the tiers after it are not asked\n", p.ID, p.Note)
} }
} }
advancePlans(ctx, open) advanceHeld(ctx, open)
} }
// advancePlans moves every open plan as far as the facts allow: a tier whose modules are all built // advancePlans moves every open plan as far as the facts allow: a tier whose modules are all built
// and whose gates are applied gives way to the next; the last tier done is the plan done. Called // and whose gates are applied gives way to the next; the last tier done is the plan done. Called
// after every outcome and on a timer, so a plan waiting on a machine's report moves when it comes. // after every outcome and on a timer, so a plan waiting on a machine's report moves when it comes.
//
// **One controller at a time** (novox/hq issue 213). A plan is read, changed and saved whole; two
// controllers — the old and the new while a machine hands its controller over — would each ask a
// tier the other had just asked. Taken without waiting: whoever holds the plans is moving them.
func advancePlans(ctx context.Context, open *stores) { func advancePlans(ctx context.Context, open *stores) {
release, err := open.inventory.HoldPlans(ctx, false)
if err != nil {
if !errors.Is(err, inventory.ErrPlansBusy) {
fmt.Printf("plans: cannot hold them: %v\n", err)
}
return
}
defer release()
advanceHeld(ctx, open)
}
// advanceHeld is advancePlans for a caller already holding the plans.
func advanceHeld(ctx context.Context, open *stores) {
inv := open.inventory inv := open.inventory
plans, err := inv.OpenPlans(ctx) plans, err := inv.OpenPlans(ctx)
if err != nil { if err != nil {
@@ -676,6 +703,19 @@ func plansCommand(ctx context.Context, args []string) error {
} }
p.State = inventory.PlanFailed p.State = inventory.PlanFailed
p.Note = "stopped by hand at tier " + fmt.Sprint(p.Tier) p.Note = "stopped by hand at tier " + fmt.Sprint(p.Tier)
release, err := inv.HoldPlans(ctx, true)
if err != nil {
return err
}
defer release()
if p, err = inv.PlanByID(ctx, positionals[1]); err != nil {
return err
}
if !p.Open() {
return fmt.Errorf("%s is already %s", p.ID, p.State)
}
p.State = inventory.PlanFailed
p.Note = "stopped by hand at tier " + fmt.Sprint(p.Tier)
if err := inv.SavePlan(ctx, p); err != nil { if err := inv.SavePlan(ctx, p); err != nil {
return err return err
} }
+7
View File
@@ -315,6 +315,13 @@ func (f following) SourceMoved(ctx context.Context, m link.SourceMoved) error {
for _, e := range moved { for _, e := range moved {
movedNames = append(movedNames, e.Manifest.Module) movedNames = append(movedNames, e.Manifest.Module)
} }
// Written and its first tier asked as one act on the plans (novox/hq issue 213): a timer on
// another controller reading it between the two would ask the tier again.
release, err := inv.HoldPlans(ctx, true)
if err != nil {
return notNow(err)
}
defer release()
plan := planOfMerge(m, movedNames, edges) plan := planOfMerge(m, movedNames, edges)
if hasCycle(plan.Tiers, edges) { if hasCycle(plan.Tiers, edges) {
fmt.Printf(" the last tier depends on itself: %s — built together, in no order\n", fmt.Printf(" the last tier depends on itself: %s — built together, in no order\n",
+59
View File
@@ -87,3 +87,62 @@ func (i *Inventory) tryHold(ctx context.Context, sorted []string) (func(), strin
} }
return release, "", nil return release, "", nil
} }
// ErrPlansBusy is the plans held by another act — on a machine replacing its controller, the other
// controller — for longer than a caller waits, or at all for one that does not wait.
var ErrPlansBusy = errors.New("another controller is working the plans")
// HoldPlans makes working the plans one act at a time, across every controller on the store
// (novox/hq issue 213). A plan is read, changed and written whole; two controllers doing that at
// once — the old and the new for the moment a machine hands its controller over, or a controller
// and a person's `plans stop` — each act on what the other has not saved yet: a tier asked twice,
// an outcome written over. A session-level advisory lock on one connection, released by the
// returned function and by the session ending, so a controller that dies holding it holds nothing.
//
// wait false gives ErrPlansBusy at once when another holds them — the timer's way: the holder is
// moving the plans already. wait true looks again every HoldPoll for up to HoldWaitFor — an
// outcome's or a merge's way, which must be written.
func (i *Inventory) HoldPlans(ctx context.Context, wait bool) (func(), error) {
deadline := time.Now().Add(HoldWaitFor)
for {
release, took, err := i.tryLock(ctx, "mesh-plans")
if err != nil || took {
return release, err
}
if !wait || time.Now().After(deadline) {
return nil, ErrPlansBusy
}
select {
case <-ctx.Done():
return nil, ctx.Err()
case <-time.After(HoldPoll):
}
}
}
// tryLock takes one named advisory lock on a connection of its own, or gives the connection back.
func (i *Inventory) tryLock(ctx context.Context, key string) (func(), bool, error) {
conn, err := i.store.Pool().Acquire(ctx)
if err != nil {
return nil, false, err
}
var once sync.Once
release := func() {
once.Do(func() {
if _, err := conn.Exec(context.WithoutCancel(ctx), `select pg_advisory_unlock_all()`); err != nil {
_ = conn.Conn().Close(context.WithoutCancel(ctx))
}
conn.Release()
})
}
var took bool
if err := conn.QueryRow(ctx, `select pg_try_advisory_lock(hashtext($1)::bigint)`, key).Scan(&took); err != nil {
release()
return nil, false, err
}
if !took {
release()
return nil, false, nil
}
return release, true, nil
}
+38
View File
@@ -0,0 +1,38 @@
package inventory
import (
"errors"
"testing"
"time"
)
// novox/hq issue 213: while a machine hands its controller over from the container to the process,
// two controllers run on one store for a moment. Working the plans is one act at a time across them.
func TestThePlansAreWorkedByOneControllerAtATime(t *testing.T) {
first := ForTest(t)
// A second controller: its own connections to the same store.
second, err := Open(t.Context())
if err != nil {
t.Fatal(err)
}
t.Cleanup(second.Close)
release, err := first.HoldPlans(t.Context(), false)
if err != nil {
t.Fatalf("the plans could not be held when nobody held them: %v", err)
}
t.Cleanup(release) // a pool closing waits for a connection still held; release is idempotent
if _, err := second.HoldPlans(t.Context(), false); !errors.Is(err, ErrPlansBusy) {
t.Fatalf("a second controller held the plans while the first did: %v", err)
}
// A waiter gets them once they are let go.
was := HoldPoll
HoldPoll = 10 * time.Millisecond
defer func() { HoldPoll = was }()
go func() { time.Sleep(50 * time.Millisecond); release() }()
again, err := second.HoldPlans(t.Context(), true)
if err != nil {
t.Fatalf("a waiting controller never got the plans once they were let go: %v", err)
}
again()
}
+54 -2
View File
@@ -5,6 +5,7 @@ import (
"encoding/json" "encoding/json"
"errors" "errors"
"fmt" "fmt"
"log"
"strings" "strings"
"time" "time"
@@ -78,10 +79,15 @@ func (n *natsInbound) Receive(ctx context.Context, act func(context.Context, Con
// than creating one here: the consumer is an object with a configuration — ack policy, ack // than creating one here: the consumer is an object with a configuration — ack policy, ack
// wait, redelivery — and a client that creates its own would be a second opinion about it. // wait, redelivery — and a client that creates its own would be a second opinion about it.
control := make(chan *nats.Msg, Prefetch) control := make(chan *nats.Msg, Prefetch)
said, err := js.ChanSubscribe("", control, nats.Bind("CONTROL", broker.ControllerName)) said, err := standingBy(ctx, log.Default(), "CONTROL", func() (*nats.Subscription, error) {
return js.ChanSubscribe("", control, nats.Bind("CONTROL", broker.ControllerName))
})
if err != nil { if err != nil {
return fmt.Errorf("subscribing to what nodes say: %w", err) return fmt.Errorf("subscribing to what nodes say: %w", err)
} }
if said == nil {
return nil // stopped while standing by
}
defer func() { _ = said.Unsubscribe() }() defer func() { _ = said.Unsubscribe() }()
// Heartbeats, on core NATS and off any stream (design 25 §3). Their own subscription because // Heartbeats, on core NATS and off any stream (design 25 §3). Their own subscription because
@@ -97,10 +103,15 @@ func (n *natsInbound) Receive(ctx context.Context, act func(context.Context, Con
var events chan *nats.Msg var events chan *nats.Msg
if len(n.follows) > 0 { if len(n.follows) > 0 {
events = make(chan *nats.Msg, Prefetch) events = make(chan *nats.Msg, Prefetch)
followed, err := js.ChanSubscribe("", events, nats.Bind("EVENTS", broker.ControllerName)) followed, err := standingBy(ctx, log.Default(), "EVENTS", func() (*nats.Subscription, error) {
return js.ChanSubscribe("", events, nats.Bind("EVENTS", broker.ControllerName))
})
if err != nil { if err != nil {
return fmt.Errorf("subscribing to what the catalogue says: %w", err) return fmt.Errorf("subscribing to what the catalogue says: %w", err)
} }
if followed == nil {
return nil
}
defer func() { _ = followed.Unsubscribe() }() defer func() { _ = followed.Unsubscribe() }()
} }
@@ -324,3 +335,44 @@ func (m *natsControl) forget() {
type replyAddressed struct { type replyAddressed struct {
ReplyTo string `json:"reply_to,omitempty"` ReplyTo string `json:"reply_to,omitempty"`
} }
// StandbyPoll is how often a controller standing by looks again for its consumers. A variable so a
// test need not wait.
var StandbyPoll = 2 * time.Second
// standingBy binds one of the controller's consumers, waiting while another controller holds it.
//
// **Two controllers, one consumer** (novox/hq issue 213). The controller's consumers are push
// consumers with no delivery group, so the server lets one subscription bind each — on purpose:
// two would each act on every message (issue 146). When a machine hands its controller over from
// the container to the process, the host starts the process first and removes the container only
// once the process is up; the process then finds the consumers bound. Exiting on that would never
// be up, so the container would never go. It stands by instead — the seat's verbs are already
// served from a queue group, and the plans wait on their lock — and binds as soon as the other lets
// go. Nil and no error is ctx ending while it waited.
func standingBy(ctx context.Context, logger interface{ Printf(string, ...any) }, stream string,
bind func() (*nats.Subscription, error)) (*nats.Subscription, error) {
said := false
for {
sub, err := bind()
if err == nil {
if said {
logger.Printf("took the controller's consumer on %s: the controller that held it let go", stream)
}
return sub, nil
}
if !strings.Contains(err.Error(), "already bound") {
return nil, err
}
if !said {
logger.Printf("another controller holds the controller's consumer on %s; standing by "+
"until it lets go", stream)
said = true
}
select {
case <-ctx.Done():
return nil, nil
case <-time.After(StandbyPoll):
}
}
}
+69
View File
@@ -0,0 +1,69 @@
package link
import (
"context"
"encoding/json"
"os"
"testing"
"time"
"github.com/novox/mesh-controller/internal/broker"
)
// novox/hq issue 213: while a machine hands its controller over, the new controller (the process)
// is started while the old one (the container) still holds the controller's consumers. It must not
// exit — the host would read that as a replacement that did not come up and never remove the
// container — and must not act on what the old one is handed. It stands by, and takes the consumers
// when the old one lets go.
func TestNatsASecondControllerStandsByAndTakesOverWhenTheFirstLetsGo(t *testing.T) {
js := aBus(t)
was := StandbyPoll
StandbyPoll = 50 * time.Millisecond
defer func() { StandbyPoll = was }()
old := &counted{}
_, stopOld := servingOn(t, js, old)
eventually(t, "the first controller binding its consumer", func() bool {
info, err := js.Context().ConsumerInfo("CONTROL", broker.ControllerName)
return err == nil && info.PushBound
})
// The new one, on a connection of its own as the process would have.
second, err := broker.Dial(os.Getenv("MESH_TEST_NATS"))
if err != nil {
t.Fatal(err)
}
t.Cleanup(second.Close)
fresh := &counted{}
s := &Server{inbound: Nats(second), bus: OverNATS{Conn: second.Conn(), JS: second.Context()},
listener: fresh, log: quiet()}
ctx, stopNew := context.WithCancel(context.Background())
defer stopNew()
ended := make(chan error, 1)
go func() { ended <- s.Serve(ctx) }()
select {
case err := <-ended:
t.Fatalf("the second controller stopped instead of standing by: %v", err)
case <-time.After(500 * time.Millisecond):
}
report := func(declared string) {
body, _ := json.Marshal(Report{Node: "anchor", Declared: declared, Applied: []string{"store"}})
if _, err := js.Context().Publish(ReportSubject("anchor"), body); err != nil {
t.Fatal(err)
}
}
report("d1")
eventually(t, "the holding controller hearing the report", func() bool { return old.count() == 1 })
if fresh.count() != 0 {
t.Fatal("the controller standing by acted on a report the holder was handed")
}
stopOld()
report("d2")
eventually(t, "the second controller taking over once the first let go", func() bool { return fresh.count() == 1 })
if old.count() != 1 {
t.Errorf("the first controller heard %d reports", old.count())
}
}