Compare commits
1
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
00bcdccdc7 |
@@ -12,6 +12,7 @@ import (
|
||||
|
||||
"git.novox.be/novox/mesh-sdk/go/asks"
|
||||
|
||||
"github.com/novox/mesh-controller/internal/broker"
|
||||
"github.com/novox/mesh-controller/internal/conditions"
|
||||
"github.com/novox/mesh-controller/internal/link"
|
||||
)
|
||||
@@ -703,3 +704,97 @@ func TestNothingIsAskedWhileTheBusLacksTheControllersGrant(t *testing.T) {
|
||||
t.Errorf("the undelivered condition was not cleared: %+v", last)
|
||||
}
|
||||
}
|
||||
|
||||
// hq issue 353, the review: the bus holds the controller's grant to ask only when both facts hold — the mesh
|
||||
// composes the controller's own grant, and the bus's module says it is in force on the running server — and
|
||||
// each fact not known is not held, said with why.
|
||||
func TestTheGrantIsHeldOnlyWhenComposedAndInForceOnTheBus(t *testing.T) {
|
||||
yes := func(context.Context) (bool, string, error) { return true, "", nil }
|
||||
no := func(why string) grantJudge {
|
||||
return func(context.Context) (bool, string, error) { return false, why, nil }
|
||||
}
|
||||
broken := func(context.Context) (bool, string, error) { return false, "", errors.New("the store is away") }
|
||||
for _, c := range []struct {
|
||||
name string
|
||||
composed, inForce grantJudge
|
||||
held bool
|
||||
says string
|
||||
}{
|
||||
{"composed and in force", yes, yes, true, ""},
|
||||
{"not composed", no("no router is assigned"), yes, false, "no router is assigned"},
|
||||
{"composed, written, not reloaded", yes, no("the server has not reloaded"), false, "has not reloaded"},
|
||||
{"the composition unknown", broken, yes, false, "could not be worked out: the store is away"},
|
||||
{"the bus unknown", yes, broken, false, "could not be worked out: the store is away"},
|
||||
{"nothing judges the bus", yes, nil, false, "not judged"},
|
||||
} {
|
||||
held, why, err := grantHeldBy(c.composed, c.inForce, time.Now)(context.Background())
|
||||
if err != nil || held != c.held || !strings.Contains(why, c.says) {
|
||||
t.Errorf("%s: held %v (%q, %v)", c.name, held, why, err)
|
||||
}
|
||||
}
|
||||
// Remembered for grantLookEvery, then judged again.
|
||||
now := time.Date(2026, 10, 9, 18, 0, 0, 0, time.UTC)
|
||||
inForce := false
|
||||
judge := grantHeldBy(yes, func(context.Context) (bool, string, error) { return inForce, "not yet", nil },
|
||||
func() time.Time { return now })
|
||||
if held, _, _ := judge(context.Background()); held {
|
||||
t.Fatal("held before the bus enforces it")
|
||||
}
|
||||
inForce = true
|
||||
if held, _, _ := judge(context.Background()); held {
|
||||
t.Error("judged again inside grantLookEvery")
|
||||
}
|
||||
now = now.Add(grantLookEvery)
|
||||
if held, _, _ := judge(context.Background()); !held {
|
||||
t.Error("not judged again after grantLookEvery")
|
||||
}
|
||||
}
|
||||
|
||||
// The controller's own grant, as composed: present while a router that takes asks is assigned, absent
|
||||
// otherwise — whatever else in the user list changes.
|
||||
func TestTheControllersOwnGrantToAskIsReadFromTheComposition(t *testing.T) {
|
||||
router := broker.Declared{Module: "messenger", Holds: []broker.Seat{{Name: broker.AsksSeat, Scope: "mesh",
|
||||
Accepts: []string{"ask", "cancel"}, Emits: []string{"decided"}, ByCaller: []string{"ask", "cancel", "decided"}}}}
|
||||
other := broker.Declared{Module: "dunst"}
|
||||
with := broker.Records{Nodes: []string{"anchor"}, Assigned: map[string][]broker.Declared{"anchor": {router, other}}}
|
||||
without := broker.Records{Nodes: []string{"anchor"}, Assigned: map[string][]broker.Declared{"anchor": {other}}}
|
||||
if may, err := controllerMayAsk(with); err != nil || !may {
|
||||
t.Errorf("with a router assigned: %v %v", may, err)
|
||||
}
|
||||
if may, err := controllerMayAsk(without); err != nil || may {
|
||||
t.Errorf("with no router assigned: %v %v", may, err)
|
||||
}
|
||||
}
|
||||
|
||||
// The bus module's answer: in force only when it says so in so many words.
|
||||
func TestTheBusModulesAnswerIsInForceOnlyWhenItSaysSo(t *testing.T) {
|
||||
for _, c := range []struct {
|
||||
answer string
|
||||
want bool
|
||||
says string
|
||||
}{
|
||||
{`{"allowed":true,"in_force":true}`, true, ""},
|
||||
{`{"allowed":true,"in_force":false,"why":"the server has not reloaded since it was written"}`, false, "not reloaded"},
|
||||
{`{"allowed":false,"in_force":false}`, false, "not in force"},
|
||||
{`{"allowed":true}`, false, "older than that answer"},
|
||||
{`{"allowed":false,"in_force":true}`, false, "not allowed"},
|
||||
{`not json`, false, "could not be read"},
|
||||
} {
|
||||
if got, why := grantInForce(json.RawMessage(c.answer)); got != c.want || !strings.Contains(why, c.says) {
|
||||
t.Errorf("%s: %v %q", c.answer, got, why)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// The serving asker is wired to the real judgement: with records and a bus it cannot read, the grant is not
|
||||
// held, and so nothing is asked (a stub saying yes, or no judgement at all, would ask).
|
||||
func TestTheServingAskerJudgesTheGrantForReal(t *testing.T) {
|
||||
a := servingAsker(&stores{}, nil, nil, nil)
|
||||
if a.grantHeld == nil {
|
||||
t.Fatal("the serving asker does not judge the bus's grant")
|
||||
}
|
||||
held, why, err := a.grantHeld(context.Background())
|
||||
if err != nil || held || !strings.Contains(why, "could not be worked out") {
|
||||
t.Errorf("held %v (%q, %v) with nothing to read", held, why, err)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -235,12 +235,27 @@ func routerHereIn(inv *inventory.Inventory) func(ctx context.Context) (bool, err
|
||||
}
|
||||
}
|
||||
|
||||
// grantHeldIn says whether the bus holds the controller's grant to ask (novox/hq issue 353): the user list the
|
||||
// machine holding the bus was last sent is the one the mesh composes now (brokerBehind, the same judgement a
|
||||
// push makes to send that machine first). While it is behind, the controller's ask is refused by the bus,
|
||||
// whatever the record says of the router, so nothing is asked and the operator is told to push that machine.
|
||||
// Judged at most every grantLookEvery: composing the list resolves the bus's machine whole.
|
||||
func grantHeldIn(open *stores) func(ctx context.Context) (bool, string, error) {
|
||||
// grantHeldIn says whether the bus holds the controller's grant to ask (novox/hq issue 353), judged on the two
|
||||
// facts that make it so, and failing closed on each:
|
||||
//
|
||||
// 1. **composed**: the user list the mesh composes now grants the controller itself to publish its ask
|
||||
// (controllerMayAsk) — the controller's own grant, not whether any user of the list changed;
|
||||
// 2. **in force**: the bus's own module, on the machine holding the bus, says the user list there allows it
|
||||
// and the server reloaded after that list was written (the nats module's `nats_user_can`, `in_force`) — a
|
||||
// list sent is not a list the bus enforces until it reloads.
|
||||
//
|
||||
// Anything that cannot be worked out is not held, with why, so the conditions that need the operator are said
|
||||
// as undelivered rather than asked into a refusal.
|
||||
func grantHeldIn(open *stores, conn *nats.Conn) func(ctx context.Context) (bool, string, error) {
|
||||
return grantHeldBy(composedGrantIn(open), busGrantIn(open, conn), time.Now)
|
||||
}
|
||||
|
||||
// grantJudge is one of the two facts: whether it holds, and why not.
|
||||
type grantJudge func(ctx context.Context) (bool, string, error)
|
||||
|
||||
// grantHeldBy is the judgement over the two facts, remembered for grantLookEvery: composing the list resolves
|
||||
// the bus's machine whole, and asking the bus's module is a call.
|
||||
func grantHeldBy(composed, inForce grantJudge, now func() time.Time) func(ctx context.Context) (bool, string, error) {
|
||||
var mu sync.Mutex
|
||||
var at time.Time
|
||||
var held bool
|
||||
@@ -248,22 +263,132 @@ func grantHeldIn(open *stores) func(ctx context.Context) (bool, string, error) {
|
||||
return func(ctx context.Context) (bool, string, error) {
|
||||
mu.Lock()
|
||||
defer mu.Unlock()
|
||||
if !at.IsZero() && time.Since(at) < grantLookEvery {
|
||||
if !at.IsZero() && now().Sub(at) < grantLookEvery {
|
||||
return held, why, nil
|
||||
}
|
||||
machine, behind, err := brokerBehind(ctx, open, nil)
|
||||
held, why = judgeGrant(ctx, composed, inForce)
|
||||
at = now()
|
||||
return held, why, nil
|
||||
}
|
||||
}
|
||||
|
||||
// judgeGrant is held only when both facts hold; an error is a fact not known, never a yes.
|
||||
func judgeGrant(ctx context.Context, composed, inForce grantJudge) (bool, string) {
|
||||
for _, fact := range []struct {
|
||||
judge grantJudge
|
||||
what string
|
||||
}{{composed, "whether the mesh composes the controller a grant to ask"}, {inForce, "whether the bus enforces it"}} {
|
||||
if fact.judge == nil {
|
||||
return false, fact.what + " is not judged here"
|
||||
}
|
||||
ok, why, err := fact.judge(ctx)
|
||||
if err != nil {
|
||||
return false, fact.what + " could not be worked out: " + err.Error()
|
||||
}
|
||||
if !ok {
|
||||
return false, why
|
||||
}
|
||||
}
|
||||
return true, ""
|
||||
}
|
||||
|
||||
// controllerMayAsk says whether the controller's own grant, as composed from these records, lets it publish its
|
||||
// ask (ADR 0259 §3: composed while a router that takes asks under its asker's name is assigned).
|
||||
func controllerMayAsk(records broker.Records) (bool, error) {
|
||||
users, err := broker.Users(records)
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
for _, u := range users {
|
||||
if u.Kind != broker.KindController {
|
||||
continue
|
||||
}
|
||||
perms, err := broker.PermissionsFor(u)
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
return broker.MayPublish(perms, asks.AskSubject(askerName)), nil
|
||||
}
|
||||
return false, errors.New("the composition holds no controller")
|
||||
}
|
||||
|
||||
// composedGrantIn is the first fact, from the mesh's records.
|
||||
func composedGrantIn(open *stores) grantJudge {
|
||||
return func(ctx context.Context) (bool, string, error) {
|
||||
if open == nil || open.inventory == nil {
|
||||
return false, "", errors.New("the mesh's records are not open")
|
||||
}
|
||||
records, err := open.inventory.BusRecords(ctx)
|
||||
if err != nil {
|
||||
return false, "", err
|
||||
}
|
||||
at, held, why = time.Now(), !behind, ""
|
||||
if behind {
|
||||
why = fmt.Sprintf("the bus's user list on %s is behind what the mesh composes, so the bus has not been "+
|
||||
"given the controller's grant to ask; `push %s` carries it", machine, machine)
|
||||
may, err := controllerMayAsk(records)
|
||||
if err != nil || may {
|
||||
return may, "", err
|
||||
}
|
||||
return held, why, nil
|
||||
return false, "the mesh composes the controller no grant to ask: no router that takes asks under its " +
|
||||
"asker's name is assigned", nil
|
||||
}
|
||||
}
|
||||
|
||||
// busGrantIn is the second fact, asked of the bus's own module on the machine holding the bus.
|
||||
func busGrantIn(open *stores, conn *nats.Conn) grantJudge {
|
||||
return func(ctx context.Context) (bool, string, error) {
|
||||
if open == nil || open.inventory == nil || conn == nil {
|
||||
return false, "", errors.New("the mesh's records or the bus are not open")
|
||||
}
|
||||
holders, err := seatHolders(ctx, open.inventory)
|
||||
if err != nil {
|
||||
return false, "", err
|
||||
}
|
||||
h, held := holders[theBrokerSeat]
|
||||
if !held || h.Node == "" || h.Module == "" {
|
||||
return false, "no module holds the bus, so nothing says what it enforces", nil
|
||||
}
|
||||
a, err := link.AskModuleToolOn(ctx, conn, h.Module, busUserCanTool, h.Node, map[string]any{
|
||||
"user": broker.ControllerName, "action": "publish", "subject": asks.AskSubject(askerName)}, 15*time.Second)
|
||||
if err != nil {
|
||||
return false, "", err
|
||||
}
|
||||
if a.Error != "" {
|
||||
return false, "", errors.New(a.Error)
|
||||
}
|
||||
inForce, why := grantInForce(a.Result)
|
||||
if !inForce {
|
||||
why = fmt.Sprintf("the bus on %s does not enforce the controller's grant to ask yet: %s; `push %s` "+
|
||||
"carries it, and the bus reloads it in place", h.Node, why, h.Node)
|
||||
}
|
||||
return inForce, why, nil
|
||||
}
|
||||
}
|
||||
|
||||
// busUserCanTool is the bus module's answer to whether a user may do something, and whether it is in force.
|
||||
const busUserCanTool = "nats_user_can"
|
||||
|
||||
// grantInForce reads the bus module's answer: in force only when it says so, in so many words. An answer
|
||||
// without `in_force` is a bus module older than the question, and is not a yes.
|
||||
func grantInForce(raw json.RawMessage) (bool, string) {
|
||||
var answer struct {
|
||||
Allowed bool `json:"allowed"`
|
||||
InForce *bool `json:"in_force"`
|
||||
Why string `json:"why"`
|
||||
}
|
||||
if err := json.Unmarshal(raw, &answer); err != nil {
|
||||
return false, "the bus module's answer could not be read"
|
||||
}
|
||||
switch {
|
||||
case answer.InForce == nil:
|
||||
return false, "the bus module does not say whether a grant is in force (it is older than that answer)"
|
||||
case !*answer.InForce && answer.Why != "":
|
||||
return false, answer.Why
|
||||
case !*answer.InForce:
|
||||
return false, "the bus module says it is not in force"
|
||||
case !answer.Allowed:
|
||||
return false, "the bus module says it is in force and not allowed"
|
||||
}
|
||||
return true, ""
|
||||
}
|
||||
|
||||
// grantLookEvery is how often the bus's user list is judged against the one its machine was last sent.
|
||||
const grantLookEvery = 30 * time.Second
|
||||
|
||||
@@ -302,7 +427,19 @@ func startAsking(ctx context.Context, open *stores, server *link.Server, conn *n
|
||||
fmt.Printf("the operator cannot be asked: %v\n", err)
|
||||
return
|
||||
}
|
||||
a := &asker{
|
||||
a := servingAsker(open, conn, js, keeper)
|
||||
if err := server.Decides(a); err != nil {
|
||||
fmt.Printf("the operator's answers cannot be heard, so nothing is asked: %v\n", err)
|
||||
return
|
||||
}
|
||||
askerFrom = a
|
||||
go a.keep(ctx)
|
||||
}
|
||||
|
||||
// servingAsker is the serving controller's asker, wired to the mesh: what it reads, publishes, performs and
|
||||
// records, and the judgements it asks by — the bus's grant to ask among them (novox/hq issue 353).
|
||||
func servingAsker(open *stores, conn *nats.Conn, js jetstream.JetStream, keeper *conditions.Keeper) *asker {
|
||||
return &asker{
|
||||
open: keeper.Open,
|
||||
silence: func(ctx context.Context, key string, d time.Duration, by, why string) error {
|
||||
_, err := keeper.Silence(ctx, key, d, by, why)
|
||||
@@ -320,7 +457,7 @@ func startAsking(ctx context.Context, open *stores, server *link.Server, conn *n
|
||||
},
|
||||
routerRecord: routerRecordOf(conn, open.inventory),
|
||||
routerHere: routerHereIn(open.inventory),
|
||||
grantHeld: grantHeldIn(open),
|
||||
grantHeld: grantHeldIn(open, conn),
|
||||
channels: channelsIn(open.inventory),
|
||||
raise: func(ctx context.Context, obs []conditions.Observation) error {
|
||||
return keeper.Reconcile(ctx, sourceAsker, obs)
|
||||
@@ -328,10 +465,4 @@ func startAsking(ctx context.Context, open *stores, server *link.Server, conn *n
|
||||
now: time.Now,
|
||||
logf: func(format string, args ...any) { fmt.Printf(format+"\n", args...) },
|
||||
}
|
||||
if err := server.Decides(a); err != nil {
|
||||
fmt.Printf("the operator's answers cannot be heard, so nothing is asked: %v\n", err)
|
||||
return
|
||||
}
|
||||
askerFrom = a
|
||||
go a.keep(ctx)
|
||||
}
|
||||
|
||||
@@ -119,6 +119,13 @@ func TestTheBusMachineIsBehindByItsUserListAlone(t *testing.T) {
|
||||
if userListBehind("", "") {
|
||||
t.Fatal("a machine sent no list reads as behind")
|
||||
}
|
||||
// A list that could not be composed is behind: not known is never "not behind" (the review of hq issue 353).
|
||||
if !listBehind("", errors.New("the plan cannot be worked out"), digestOf([]byte(list))) {
|
||||
t.Fatal("a list nobody could compose reads as current")
|
||||
}
|
||||
if listBehind(list, nil, digestOf([]byte(list))) || !listBehind(list, nil, "") {
|
||||
t.Fatal("listBehind does not read the digest as userListBehind does")
|
||||
}
|
||||
}
|
||||
|
||||
// The machine holding the bus goes first: its declaration carries the user list the new grants are
|
||||
|
||||
@@ -1096,20 +1096,29 @@ func brokerBehind(ctx context.Context, open *stores, names []string) (string, bo
|
||||
return h.Node, false, nil
|
||||
}
|
||||
}
|
||||
// **What cannot be worked out is behind** (the review of hq issue 353). The list it would be sent cannot be
|
||||
// composed, so whether it carries the grants a send relies on is not known; it is taken as behind, so a push
|
||||
// sends that machine first and says why it cannot (`plan` says it too), and a rollout that needs the list
|
||||
// refuses rather than sending code whose grants the bus may not hold. It said "not behind" until then: a
|
||||
// send went ahead as if the bus held a list nobody could compose.
|
||||
list, composeErr := "", error(nil)
|
||||
plan, _, err := planFor(ctx, open, h.Node)
|
||||
if err != nil {
|
||||
// It cannot be worked out: sending it would refuse the whole send, and `plan` says why.
|
||||
return h.Node, false, nil
|
||||
}
|
||||
list, _, err := busUserList(ctx, inv, plan.Modules)
|
||||
if err != nil {
|
||||
return h.Node, false, nil
|
||||
composeErr = err
|
||||
} else if list, _, err = busUserList(ctx, inv, plan.Modules); err != nil {
|
||||
composeErr = err
|
||||
}
|
||||
sent, err := inv.SentBusUsers(ctx, h.Node)
|
||||
if err != nil {
|
||||
return "", false, err
|
||||
}
|
||||
return h.Node, userListBehind(list, sent), nil
|
||||
return h.Node, listBehind(list, composeErr, sent), nil
|
||||
}
|
||||
|
||||
// listBehind is whether the bus's machine is behind: the list composed now is not the one last sent, or it
|
||||
// could not be composed at all — not known is behind, never "not behind".
|
||||
func listBehind(list string, composeErr error, sentDigest string) bool {
|
||||
return composeErr != nil || userListBehind(list, sentDigest)
|
||||
}
|
||||
|
||||
// userListBehind is whether the user list composed now is not the one last sent, by its digest. An
|
||||
|
||||
Reference in New Issue
Block a user