Merge pull request 'Ask the operator only once the bus holds the controller's grant to ask (hq issue 353)' (#186) from fix/353-the-controller-asks-only-once-the-bus-holds-its-grant into main
This commit was merged in pull request #186.
This commit is contained in:
@@ -167,6 +167,11 @@ type asker struct {
|
|||||||
routerRecord func(ctx context.Context, id string) (*asks.Warrant, error)
|
routerRecord func(ctx context.Context, id string) (*asks.Warrant, error)
|
||||||
// routerHere says whether a router holds the seat and takes asks under the asker's name; nil is yes.
|
// routerHere says whether a router holds the seat and takes asks under the asker's name; nil is yes.
|
||||||
routerHere func(ctx context.Context) (bool, error)
|
routerHere func(ctx context.Context) (bool, error)
|
||||||
|
// grantHeld says whether the bus holds the controller's grant to ask: the user list the bus's machine was
|
||||||
|
// last sent is the one the mesh composes now (novox/hq issue 353). The grant is composed from the router's
|
||||||
|
// assignment and reaches the bus only when that machine is next pushed, so between `assign` and `push`
|
||||||
|
// the record says a router is here and the bus refuses every ask. why says what to do; nil is yes.
|
||||||
|
grantHeld func(ctx context.Context) (held bool, why string, err error)
|
||||||
// channels is what the channels are now, as a fingerprint: who holds which kind, promising what.
|
// channels is what the channels are now, as a fingerprint: who holds which kind, promising what.
|
||||||
channels func(ctx context.Context) string
|
channels func(ctx context.Context) string
|
||||||
// raise keeps the asker's own condition (sourceAsker): which conditions needing the operator could not be
|
// raise keeps the asker's own condition (sourceAsker): which conditions needing the operator could not be
|
||||||
@@ -176,6 +181,7 @@ type asker struct {
|
|||||||
logf func(string, ...any)
|
logf func(string, ...any)
|
||||||
|
|
||||||
saidNoRouter bool
|
saidNoRouter bool
|
||||||
|
saidNoGrant bool
|
||||||
|
|
||||||
mu sync.Mutex
|
mu sync.Mutex
|
||||||
nudged chan struct{}
|
nudged chan struct{}
|
||||||
@@ -258,6 +264,30 @@ func (a *asker) reconcile(ctx context.Context) error {
|
|||||||
}
|
}
|
||||||
a.saidNoRouter = false
|
a.saidNoRouter = false
|
||||||
}
|
}
|
||||||
|
if a.grantHeld != nil {
|
||||||
|
held, why, err := a.grantHeld(ctx)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
if !held {
|
||||||
|
if !a.saidNoGrant {
|
||||||
|
a.logf("the bus does not hold the controller's grant to ask yet: %s; the operator is asked nothing until it does", why)
|
||||||
|
a.saidNoGrant = true
|
||||||
|
}
|
||||||
|
open, err := a.open(ctx)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
var unasked []conditions.Condition
|
||||||
|
for _, c := range open {
|
||||||
|
if wants(c, now) {
|
||||||
|
unasked = append(unasked, c)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return a.sayUnasked(ctx, unasked, "the bus does not hold the controller's grant to ask yet: "+why)
|
||||||
|
}
|
||||||
|
a.saidNoGrant = false
|
||||||
|
}
|
||||||
channels := ""
|
channels := ""
|
||||||
if a.channels != nil {
|
if a.channels != nil {
|
||||||
channels = a.channels(ctx)
|
channels = a.channels(ctx)
|
||||||
|
|||||||
@@ -4,6 +4,7 @@ import (
|
|||||||
"context"
|
"context"
|
||||||
"encoding/json"
|
"encoding/json"
|
||||||
"errors"
|
"errors"
|
||||||
|
"fmt"
|
||||||
"strings"
|
"strings"
|
||||||
"sync"
|
"sync"
|
||||||
"testing"
|
"testing"
|
||||||
@@ -654,3 +655,51 @@ func TestAnAcknowledgementNeverSharesAnAskWithAnApproval(t *testing.T) {
|
|||||||
t.Errorf("the approval kept through a silence was not performed: %v", r.called)
|
t.Errorf("the approval kept through a silence was not performed: %v", r.called)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// novox/hq issue 353: the controller's grant to ask is composed from the router's assignment and reaches the
|
||||||
|
// bus only when its machine is pushed. Between the two, the bus refuses every ask (measured 2026-10-09, 17:54 to
|
||||||
|
// 17:56 local: seven refusals of mesh.seat.operator-channel.accept.ask.mesh-controller). So nothing is asked
|
||||||
|
// while the bus's user list is behind, it is said once, the conditions that need the operator are raised as
|
||||||
|
// undelivered with what to do, and the asks go out once the bus holds the grant.
|
||||||
|
func TestNothingIsAskedWhileTheBusLacksTheControllersGrant(t *testing.T) {
|
||||||
|
r := newAskerRig(t)
|
||||||
|
var said []string
|
||||||
|
r.a.logf = func(f string, a ...any) { said = append(said, fmt.Sprintf(f, a...)) }
|
||||||
|
var raised [][]conditions.Observation
|
||||||
|
r.a.raise = func(_ context.Context, obs []conditions.Observation) error {
|
||||||
|
raised = append(raised, obs)
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
held := false
|
||||||
|
r.a.grantHeld = func(context.Context) (bool, string, error) {
|
||||||
|
return held, "the bus's user list on anchor is behind what the mesh composes; `push anchor` carries it", nil
|
||||||
|
}
|
||||||
|
r.open = []conditions.Condition{heldCondition()}
|
||||||
|
_ = r.a.reconcile(context.Background())
|
||||||
|
_ = r.a.reconcile(context.Background())
|
||||||
|
if len(r.asksSent(t)) != 0 {
|
||||||
|
t.Error("asked while the bus lacks the grant")
|
||||||
|
}
|
||||||
|
n := 0
|
||||||
|
for _, s := range said {
|
||||||
|
if strings.Contains(s, "does not hold the controller's grant") {
|
||||||
|
n++
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if n != 1 {
|
||||||
|
t.Errorf("said %d times: %q", n, said)
|
||||||
|
}
|
||||||
|
if len(raised) == 0 || len(raised[len(raised)-1]) != 1 ||
|
||||||
|
!strings.Contains(raised[len(raised)-1][0].Summary, "`push anchor` carries it") ||
|
||||||
|
raised[len(raised)-1][0].Kind != "asks-undelivered" {
|
||||||
|
t.Fatalf("not said as a condition with what to do: %+v", raised)
|
||||||
|
}
|
||||||
|
held = true
|
||||||
|
_ = r.a.reconcile(context.Background())
|
||||||
|
if sent := r.asksSent(t); len(sent) != 1 || sent[0].About != heldCondition().Key {
|
||||||
|
t.Errorf("not asked once the bus holds the grant: %+v", sent)
|
||||||
|
}
|
||||||
|
if last := raised[len(raised)-1]; len(last) != 0 {
|
||||||
|
t.Errorf("the undelivered condition was not cleared: %+v", last)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -11,6 +11,7 @@ import (
|
|||||||
"fmt"
|
"fmt"
|
||||||
"sort"
|
"sort"
|
||||||
"strings"
|
"strings"
|
||||||
|
"sync"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/nats-io/nats.go"
|
"github.com/nats-io/nats.go"
|
||||||
@@ -234,6 +235,38 @@ 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) {
|
||||||
|
var mu sync.Mutex
|
||||||
|
var at time.Time
|
||||||
|
var held bool
|
||||||
|
var why string
|
||||||
|
return func(ctx context.Context) (bool, string, error) {
|
||||||
|
mu.Lock()
|
||||||
|
defer mu.Unlock()
|
||||||
|
if !at.IsZero() && time.Since(at) < grantLookEvery {
|
||||||
|
return held, why, nil
|
||||||
|
}
|
||||||
|
machine, behind, err := brokerBehind(ctx, open, nil)
|
||||||
|
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)
|
||||||
|
}
|
||||||
|
return held, why, nil
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// grantLookEvery is how often the bus's user list is judged against the one its machine was last sent.
|
||||||
|
const grantLookEvery = 30 * time.Second
|
||||||
|
|
||||||
// channelsIn is what the channels are now, as a fingerprint: each module claiming a kind of the channel
|
// channelsIn is what the channels are now, as a fingerprint: each module claiming a kind of the channel
|
||||||
// bench, where, promising what, and whether of its own account. An ask the router refused is asked again
|
// bench, where, promising what, and whether of its own account. An ask the router refused is asked again
|
||||||
// once this changes.
|
// once this changes.
|
||||||
@@ -287,6 +320,7 @@ func startAsking(ctx context.Context, open *stores, server *link.Server, conn *n
|
|||||||
},
|
},
|
||||||
routerRecord: routerRecordOf(conn, open.inventory),
|
routerRecord: routerRecordOf(conn, open.inventory),
|
||||||
routerHere: routerHereIn(open.inventory),
|
routerHere: routerHereIn(open.inventory),
|
||||||
|
grantHeld: grantHeldIn(open),
|
||||||
channels: channelsIn(open.inventory),
|
channels: channelsIn(open.inventory),
|
||||||
raise: func(ctx context.Context, obs []conditions.Observation) error {
|
raise: func(ctx context.Context, obs []conditions.Observation) error {
|
||||||
return keeper.Reconcile(ctx, sourceAsker, obs)
|
return keeper.Reconcile(ctx, sourceAsker, obs)
|
||||||
|
|||||||
Reference in New Issue
Block a user