Answer every seat call within ten seconds and keep what came of it (hq issue 265)
A push outlasted the console's 30s wait and, when it sent the bus its changed user list, the broker's reload forgot the reply it may send: the push happened and its caller was told it did not answer. Calls now answer in full or as running with an id, a push answers before it sends, refused answers are recorded on their call, and 'calls' reads them back.
This commit is contained in:
@@ -18,6 +18,7 @@ import (
|
||||
"regexp"
|
||||
"sort"
|
||||
"strings"
|
||||
"time"
|
||||
)
|
||||
|
||||
// A Kind is what a principal is, which decides the shape of its authority rather than its
|
||||
@@ -189,6 +190,13 @@ type Permissions struct {
|
||||
AllowResponses bool
|
||||
}
|
||||
|
||||
// ResponseTTL is how long the bus lets a principal answer a request it received. Its one answer has
|
||||
// to come inside this, and a seat's holder answers within link.AnswerWithin — inside it by design.
|
||||
// **A broker reloading its user list forgets every answer it was about to permit**, whatever this
|
||||
// says (novox/hq issue 265): a call that is still running when the list reloads has its answer
|
||||
// refused, which is why a holder answers before it does what can reload it.
|
||||
const ResponseTTL = time.Minute
|
||||
|
||||
// PermissionsFor derives a principal's authority. Pure, and the only place authority is decided:
|
||||
// a permission that cannot be derived from a declaration is a permission nobody can explain.
|
||||
func PermissionsFor(p Principal) (Permissions, error) {
|
||||
@@ -778,7 +786,7 @@ func ComposeAccounts(principals []Principal) (string, error) {
|
||||
fmt.Fprintf(&b, " publish: { allow: [%s] }\n", quoted(perms.Publish))
|
||||
fmt.Fprintf(&b, " subscribe: { allow: [%s] }\n", quoted(perms.Subscribe))
|
||||
if perms.AllowResponses {
|
||||
b.WriteString(" allow_responses: { max: 1, ttl: \"1m\" }\n")
|
||||
fmt.Fprintf(&b, " allow_responses: { max: 1, ttl: \"%dm\" }\n", int(ResponseTTL/time.Minute))
|
||||
}
|
||||
b.WriteString(" } }\n")
|
||||
}
|
||||
|
||||
@@ -73,6 +73,13 @@ var ControllerVerbs = []Verb{
|
||||
{Name: "tools", Description: "Every seat's tools, from the mesh's own records: what each role " +
|
||||
"answers, whether or not its holder is up. The mesh's own verbs are the mesh-controller seat's.",
|
||||
Input: schema(nil, nil)},
|
||||
{Name: "calls", Description: "The calls this controller answered lately and what came of each — " +
|
||||
"one still running, one that finished after its caller was told it was running, one whose answer " +
|
||||
"the bus refused — and, given a call's id, its whole answer (novox/hq issue 265). A call that has " +
|
||||
"not finished within ten seconds answers that it is running, with its id; this is where it ends.",
|
||||
Input: schema(map[string]string{
|
||||
"call": "a call's id, as a running answer or `calls` gives it: that call and its whole answer",
|
||||
}, nil)},
|
||||
{Name: "status", Description: "What is wrong, what is quiet, what is out of date, and which " +
|
||||
"machines are behind what the mesh would send them.",
|
||||
Input: schema(nil, nil)},
|
||||
@@ -128,7 +135,9 @@ var ControllerVerbs = []Verb{
|
||||
{Name: "unpin", Description: "Take that choice back, putting the question to the mesh again.",
|
||||
Input: schema(map[string]string{"node": "the machine's name", "provision": "the provision"}, []string{"node", "provision"})},
|
||||
{Name: "push", Description: "Send one machine everything it should be. With no machine named it is a push of the " +
|
||||
"WHOLE mesh — every machine that is behind — and the answer says so first; behind says that outright.",
|
||||
"WHOLE mesh — every machine that is behind — and the answer says so first; behind says that outright. " +
|
||||
"Answers at once that it is running, with a call id: `calls` with that id says what it sent " +
|
||||
"(a push can reload the bus, which then refuses any answer still to come).",
|
||||
Input: schema(map[string]string{
|
||||
"node": "the machine's name; without it, every machine that is behind",
|
||||
"behind": "\"true\": every machine that is behind, the whole mesh — the same as naming none, said outright; not with node",
|
||||
|
||||
@@ -0,0 +1,329 @@
|
||||
package link
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"log"
|
||||
"regexp"
|
||||
"strconv"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"github.com/nats-io/nats.go"
|
||||
)
|
||||
|
||||
// A role's tool call, kept after it is answered (novox/hq issue 265).
|
||||
//
|
||||
// **A call that outlasts its caller must not end in silence.** A caller waits a bounded time (the
|
||||
// console thirty seconds), and the bus lets a holder answer a request exactly once and only while the
|
||||
// server still remembers it was asked (`allow_responses`). Before this, a push that ran longer than
|
||||
// its caller waited did everything it was asked and its caller read "did not answer in time"; and a
|
||||
// push whose first act sent the bus's own machine a changed user list had its answer refused
|
||||
// outright — a broker reloading its users forgets every reply it was about to permit — so the
|
||||
// answer went nowhere and the only trace was the client library's line on the controller's standard
|
||||
// error. So every call is answered within AnswerWithin, with its whole answer when it has one by
|
||||
// then and with "still running as call <id>" when it does not; and every call is kept here, with
|
||||
// what came of it, including an answer the bus refused, so `calls` can say it.
|
||||
|
||||
// AnswerWithin is how long a call runs before its caller is answered that it is still running. Well
|
||||
// inside the shortest wait of a caller the mesh ships (the console's thirty seconds) and the bus's
|
||||
// own window for an answer (broker.ResponseTTL), so the one answer a call has is never late for
|
||||
// either. A variable so a test need not wait.
|
||||
var AnswerWithin = 10 * time.Second
|
||||
|
||||
// KeptCalls is how many calls are kept, newest first; keptAnswer the largest answer kept of a call
|
||||
// whose caller was sent it. One its caller never had is kept whole.
|
||||
const (
|
||||
KeptCalls = 100
|
||||
keptAnswer = 64 << 10
|
||||
)
|
||||
|
||||
// The states a kept call is in.
|
||||
const (
|
||||
CallRunning = "running"
|
||||
CallAnswered = "answered"
|
||||
// CallFinishedAfter is a call that finished after its caller was told it was still running: its
|
||||
// answer is here and nowhere else.
|
||||
CallFinishedAfter = "finished after its caller was answered"
|
||||
)
|
||||
|
||||
// Call is one call of a role's tool, as `calls` shows it.
|
||||
type Call struct {
|
||||
ID string `json:"call"`
|
||||
Seat string `json:"seat"`
|
||||
Verb string `json:"verb"`
|
||||
Args json.RawMessage `json:"arguments,omitempty"`
|
||||
Started time.Time `json:"started"`
|
||||
Finished *time.Time `json:"finished,omitempty"`
|
||||
State string `json:"state"`
|
||||
// Failed is the call's own answer being an error — an answer, not a timeout.
|
||||
Failed bool `json:"failed,omitempty"`
|
||||
// Answer is what the call answered, or would have: kept for a call whose caller did not get it.
|
||||
Answer json.RawMessage `json:"answer,omitempty"`
|
||||
// Refused is the bus refusing the answer this holder sent: the caller got nothing, and this is
|
||||
// the only place that says what it would have.
|
||||
Refused string `json:"answer refused by the bus,omitempty"`
|
||||
|
||||
reply string // the subject the answer went to, which the bus names when it refuses it
|
||||
}
|
||||
|
||||
// CallLog keeps the latest calls a holder served.
|
||||
type CallLog struct {
|
||||
mu sync.Mutex
|
||||
calls []*Call // oldest first
|
||||
next uint64
|
||||
now func() time.Time
|
||||
// Follow is the tool that reads this log back, named in a running answer — set by a holder that
|
||||
// serves one (the controller's `calls`); without it, the answer points at the holder's journal.
|
||||
Follow string
|
||||
}
|
||||
|
||||
// Calls is this process's log: one holder process serves its seats on one connection.
|
||||
var Calls = NewCallLog()
|
||||
|
||||
func NewCallLog() *CallLog { return &CallLog{now: time.Now} }
|
||||
|
||||
func (l *CallLog) begin(seat, verb string, args json.RawMessage, reply string) *Call {
|
||||
l.mu.Lock()
|
||||
defer l.mu.Unlock()
|
||||
l.next++
|
||||
c := &Call{ID: "call-" + strconv.FormatInt(l.now().UnixNano(), 10) + "-" + strconv.FormatUint(l.next, 10),
|
||||
Seat: seat, Verb: verb, Args: args, Started: l.now(), State: CallRunning, reply: reply}
|
||||
l.calls = append(l.calls, c)
|
||||
if len(l.calls) > KeptCalls {
|
||||
l.calls = l.calls[len(l.calls)-KeptCalls:]
|
||||
}
|
||||
return c
|
||||
}
|
||||
|
||||
// finish records a call's answer; answeredAlready is its caller having been told it was running.
|
||||
func (l *CallLog) finish(c *Call, answer []byte, failed, answeredAlready bool) {
|
||||
l.mu.Lock()
|
||||
defer l.mu.Unlock()
|
||||
at := l.now()
|
||||
c.Finished = &at
|
||||
c.Failed = failed
|
||||
c.Answer = append(json.RawMessage(nil), answer...)
|
||||
if !answeredAlready && len(answer) > keptAnswer {
|
||||
// Its caller has it; kept only in case the bus refuses it, and a whole declaration a
|
||||
// hundred times over is memory nobody asked for.
|
||||
c.Answer, _ = json.Marshal(map[string]any{"not kept": fmt.Sprintf(
|
||||
"an answer of %d bytes, sent to its caller in full", len(answer))})
|
||||
}
|
||||
if answeredAlready {
|
||||
c.State = CallFinishedAfter
|
||||
} else {
|
||||
c.State = CallAnswered
|
||||
}
|
||||
}
|
||||
|
||||
// Recent is the kept calls, newest first, as copies.
|
||||
func (l *CallLog) Recent() []Call {
|
||||
l.mu.Lock()
|
||||
defer l.mu.Unlock()
|
||||
out := make([]Call, 0, len(l.calls))
|
||||
for i := len(l.calls) - 1; i >= 0; i-- {
|
||||
out = append(out, *l.calls[i])
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
// Get is one kept call.
|
||||
func (l *CallLog) Get(id string) (Call, bool) {
|
||||
l.mu.Lock()
|
||||
defer l.mu.Unlock()
|
||||
for _, c := range l.calls {
|
||||
if c.ID == id {
|
||||
return *c, true
|
||||
}
|
||||
}
|
||||
return Call{}, false
|
||||
}
|
||||
|
||||
// kept is what a call's arguments are kept as: each argument by name, a value only when it is short
|
||||
// and not one that carries settings or a secret — `calls` answers anyone who may call the seat.
|
||||
func kept(args json.RawMessage) json.RawMessage {
|
||||
var given map[string]any
|
||||
if json.Unmarshal(args, &given) != nil {
|
||||
return nil
|
||||
}
|
||||
out := map[string]any{}
|
||||
for k, v := range given {
|
||||
s, isString := v.(string)
|
||||
switch {
|
||||
case k == "values" || k == "secret":
|
||||
out[k] = "(given, not kept)"
|
||||
case isString && len(s) <= 120:
|
||||
out[k] = s
|
||||
case isString:
|
||||
out[k] = s[:120] + "…"
|
||||
default:
|
||||
out[k] = v
|
||||
}
|
||||
}
|
||||
body, _ := json.Marshal(out)
|
||||
return body
|
||||
}
|
||||
|
||||
// refusedPublish is the bus's words for a publish it refused, with the subject.
|
||||
var refusedPublish = regexp.MustCompile(`Permissions Violation for Publish to "([^"]+)"`)
|
||||
|
||||
// Refusal records the bus refusing an answer, from the error the client library hands the
|
||||
// connection's error handler. It reports whether the error was the refusal of a kept call's answer.
|
||||
func (l *CallLog) Refusal(err error, logger *log.Logger) bool {
|
||||
if err == nil {
|
||||
return false
|
||||
}
|
||||
m := refusedPublish.FindStringSubmatch(err.Error())
|
||||
if m == nil {
|
||||
return false
|
||||
}
|
||||
l.mu.Lock()
|
||||
var hit *Call
|
||||
for i := len(l.calls) - 1; i >= 0; i-- {
|
||||
if c := l.calls[i]; c.reply != "" && c.reply == m[1] {
|
||||
hit = c
|
||||
break
|
||||
}
|
||||
}
|
||||
if hit == nil {
|
||||
l.mu.Unlock()
|
||||
return false
|
||||
}
|
||||
at := l.now()
|
||||
hit.Refused = fmt.Sprintf("%s: %s", at.Format(time.RFC3339), err)
|
||||
id, verb, took := hit.ID, hit.Seat+"."+hit.Verb, at.Sub(hit.Started).Round(time.Second)
|
||||
l.mu.Unlock()
|
||||
if logger != nil {
|
||||
// Said in the mesh's words, beside the library's own line: which call, and where its answer is.
|
||||
logger.Printf("the bus refused the answer to %s (%s, asked %s ago): its caller got no answer. "+
|
||||
"What it answered is kept under %s. A broker reloading its user list while a call runs "+
|
||||
"forgets that it may be answered", id, verb, took, id)
|
||||
}
|
||||
return true
|
||||
}
|
||||
|
||||
// WatchRefusals chains the connection's error handler so a refused answer is recorded against its
|
||||
// call. The handler the dialler set — the library's default, which prints — still runs after it.
|
||||
func (l *CallLog) WatchRefusals(conn *nats.Conn, logger *log.Logger) {
|
||||
if conn == nil {
|
||||
return
|
||||
}
|
||||
watched.Lock()
|
||||
defer watched.Unlock()
|
||||
if watched.conns == nil {
|
||||
watched.conns = map[*nats.Conn]bool{}
|
||||
}
|
||||
if watched.conns[conn] {
|
||||
return
|
||||
}
|
||||
watched.conns[conn] = true
|
||||
before := conn.ErrorHandler()
|
||||
conn.SetErrorHandler(func(c *nats.Conn, s *nats.Subscription, err error) {
|
||||
if errors.Is(err, nats.ErrPermissionViolation) {
|
||||
l.Refusal(err, logger)
|
||||
}
|
||||
if before != nil {
|
||||
before(c, s, err)
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
var watched struct {
|
||||
sync.Mutex
|
||||
conns map[*nats.Conn]bool
|
||||
}
|
||||
|
||||
// answerNow is how a handler asks for its caller to be answered before it goes on (Acknowledge).
|
||||
type answerNowKey struct{}
|
||||
|
||||
// Acknowledge answers the call's caller now that it is running, rather than after AnswerWithin. For
|
||||
// a handler about to do what makes its answer unsendable: a push sends the bus's own machine first,
|
||||
// and a broker reloading its user list forgets every answer it was about to permit. Without effect
|
||||
// outside a served call, and after the first time.
|
||||
func Acknowledge(ctx context.Context) {
|
||||
if f, ok := ctx.Value(answerNowKey{}).(func()); ok {
|
||||
f()
|
||||
}
|
||||
}
|
||||
|
||||
// running is the interim answer: what a caller reads when the call has not finished.
|
||||
func running(c *Call, within time.Duration, acknowledged bool, follow string) []byte {
|
||||
why := fmt.Sprintf("it has not finished in %s", within)
|
||||
if acknowledged {
|
||||
why = "it answers before it starts, because what it does can leave the bus unable to carry a later answer"
|
||||
}
|
||||
next := "its holder's journal says how it ended."
|
||||
if follow != "" {
|
||||
next = fmt.Sprintf("%s with call %q says what came of it.", follow, c.ID)
|
||||
}
|
||||
said := fmt.Sprintf("%s.%s is running as %s — %s. This is not a failure, and it may already have "+
|
||||
"done what was asked: %s", c.Seat, c.Verb, c.ID, why, next)
|
||||
body, _ := json.Marshal(map[string]any{"result": map[string]any{
|
||||
"running": true, "call": c.ID, "started": c.Started, "output": said,
|
||||
}})
|
||||
return body
|
||||
}
|
||||
|
||||
// serveCall runs one call and answers it once: in full when the handler returns within AnswerWithin
|
||||
// and has not acknowledged, and otherwise that it is running — its answer then kept, and logged.
|
||||
func (l *CallLog) serveCall(seat, verb string, args json.RawMessage, reply string, handle ToolHandler,
|
||||
respond func([]byte) error, logger *log.Logger) {
|
||||
ctx, cancel := context.WithTimeout(context.Background(), HandlerTimeout)
|
||||
defer cancel()
|
||||
if len(args) == 0 {
|
||||
args = json.RawMessage(`{}`)
|
||||
}
|
||||
c := l.begin(seat, verb, kept(args), reply)
|
||||
acknowledged := make(chan struct{})
|
||||
var once sync.Once
|
||||
ctx = context.WithValue(ctx, answerNowKey{}, func() { once.Do(func() { close(acknowledged) }) })
|
||||
|
||||
type outcome struct {
|
||||
body []byte
|
||||
failed bool
|
||||
}
|
||||
done := make(chan outcome, 1)
|
||||
go func() {
|
||||
var body []byte
|
||||
result, err := handle(ctx, args)
|
||||
failed := err != nil
|
||||
if err != nil {
|
||||
body, _ = json.Marshal(map[string]any{"error": err.Error()})
|
||||
} else if body, err = json.Marshal(map[string]any{"result": result}); err != nil {
|
||||
failed = true
|
||||
body, _ = json.Marshal(map[string]any{"error": "the answer could not be written as JSON: " + err.Error()})
|
||||
}
|
||||
done <- outcome{body, failed}
|
||||
}()
|
||||
|
||||
say := func(body []byte) {
|
||||
if err := respond(body); err != nil && logger != nil {
|
||||
logger.Printf("%s.%s: could not answer %s: %v", seat, verb, c.ID, err)
|
||||
}
|
||||
}
|
||||
timer := time.NewTimer(AnswerWithin)
|
||||
defer timer.Stop()
|
||||
select {
|
||||
case o := <-done:
|
||||
l.finish(c, o.body, o.failed, false)
|
||||
say(o.body)
|
||||
return
|
||||
case <-acknowledged:
|
||||
say(running(c, AnswerWithin, true, l.Follow))
|
||||
case <-timer.C:
|
||||
say(running(c, AnswerWithin, false, l.Follow))
|
||||
}
|
||||
o := <-done
|
||||
l.finish(c, o.body, o.failed, true)
|
||||
if logger != nil {
|
||||
how := "and it succeeded"
|
||||
if o.failed {
|
||||
how = "and it answered an error"
|
||||
}
|
||||
logger.Printf("%s (%s.%s) finished after %s, after its caller was told it was running, %s — its answer is kept under %s", c.ID,
|
||||
seat, verb, time.Since(c.Started).Round(time.Second), how, c.ID)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,174 @@
|
||||
package link
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"log"
|
||||
"strings"
|
||||
"sync"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
// answers collects what a call answered, and fails a second answer: the bus permits one.
|
||||
type answers struct {
|
||||
t *testing.T
|
||||
mu sync.Mutex
|
||||
got [][]byte
|
||||
sent chan struct{}
|
||||
}
|
||||
|
||||
func newAnswers(t *testing.T) *answers { return &answers{t: t, sent: make(chan struct{}, 4)} }
|
||||
|
||||
func (a *answers) respond(body []byte) error {
|
||||
a.mu.Lock()
|
||||
defer a.mu.Unlock()
|
||||
a.got = append(a.got, body)
|
||||
if len(a.got) > 1 {
|
||||
a.t.Errorf("a call was answered %d times; the bus refuses every answer after the first", len(a.got))
|
||||
}
|
||||
a.sent <- struct{}{}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (a *answers) only() map[string]any {
|
||||
a.mu.Lock()
|
||||
defer a.mu.Unlock()
|
||||
if len(a.got) != 1 {
|
||||
a.t.Fatalf("%d answers, want exactly one", len(a.got))
|
||||
}
|
||||
var out map[string]any
|
||||
if err := json.Unmarshal(a.got[0], &out); err != nil {
|
||||
a.t.Fatal(err)
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
func shortWindow(t *testing.T, d time.Duration) {
|
||||
t.Helper()
|
||||
was := AnswerWithin
|
||||
AnswerWithin = d
|
||||
t.Cleanup(func() { AnswerWithin = was })
|
||||
}
|
||||
|
||||
// A call that finishes in time answers in full, once, and is kept as answered.
|
||||
func TestACallThatFinishesInTimeAnswersInFull(t *testing.T) {
|
||||
l, a := NewCallLog(), newAnswers(t)
|
||||
l.serveCall("mesh-controller", "status", nil, "_INBOX.x.1", func(context.Context, json.RawMessage) (any, error) {
|
||||
return "all well", nil
|
||||
}, a.respond, nil)
|
||||
if got := a.only()["result"]; got != "all well" {
|
||||
t.Fatalf("answered %v", got)
|
||||
}
|
||||
recent := l.Recent()
|
||||
if len(recent) != 1 || recent[0].State != CallAnswered || string(recent[0].Args) != "{}" {
|
||||
t.Fatalf("kept %+v", recent)
|
||||
}
|
||||
}
|
||||
|
||||
// **A call that outlasts AnswerWithin is answered that it is running, with its id** — before its
|
||||
// caller gives up — and what it finally answered is kept under that id, not dropped (issue 265).
|
||||
func TestACallThatOutlastsTheWindowSaysItIsRunningAndKeepsItsAnswer(t *testing.T) {
|
||||
shortWindow(t, 20*time.Millisecond)
|
||||
l, a := NewCallLog(), newAnswers(t)
|
||||
var logged bytes.Buffer
|
||||
release := make(chan struct{})
|
||||
finished := make(chan struct{})
|
||||
go func() {
|
||||
l.serveCall("mesh-controller", "push", json.RawMessage(`{"node":"anchor"}`), "_INBOX.x.2",
|
||||
func(context.Context, json.RawMessage) (any, error) {
|
||||
<-release
|
||||
return "anchor told", nil
|
||||
}, a.respond, log.New(&logged, "", 0))
|
||||
close(finished)
|
||||
}()
|
||||
<-a.sent
|
||||
got := a.only()["result"].(map[string]any)
|
||||
id, _ := got["call"].(string)
|
||||
if got["running"] != true || id == "" || !strings.Contains(got["output"].(string), id) {
|
||||
t.Fatalf("the running answer does not name its call: %v", got)
|
||||
}
|
||||
if c, _ := l.Get(id); c.State != CallRunning {
|
||||
t.Fatalf("while it runs it is kept as %q", c.State)
|
||||
}
|
||||
close(release)
|
||||
<-finished
|
||||
c, ok := l.Get(id)
|
||||
if !ok || c.State != CallFinishedAfter || !strings.Contains(string(c.Answer), "anchor told") {
|
||||
t.Fatalf("its answer was not kept: %+v", c)
|
||||
}
|
||||
if !strings.Contains(logged.String(), id) {
|
||||
t.Errorf("finishing late was not said: %q", logged.String())
|
||||
}
|
||||
}
|
||||
|
||||
// A handler that acknowledges is answered then, not when it ends: a push answers before it sends.
|
||||
func TestAnAcknowledgedCallIsAnsweredBeforeItGoesOn(t *testing.T) {
|
||||
l, a := NewCallLog(), newAnswers(t)
|
||||
done := make(chan struct{})
|
||||
go func() {
|
||||
l.serveCall("mesh-controller", "push", nil, "_INBOX.x.3", func(ctx context.Context, _ json.RawMessage) (any, error) {
|
||||
Acknowledge(ctx)
|
||||
Acknowledge(ctx) // a second time is nothing
|
||||
select {
|
||||
case <-a.sent: // the caller was answered before the work goes on
|
||||
case <-time.After(5 * time.Second):
|
||||
t.Error("acknowledging did not answer the caller")
|
||||
}
|
||||
return "sent", nil
|
||||
}, a.respond, nil)
|
||||
close(done)
|
||||
}()
|
||||
<-done
|
||||
if got := a.only()["result"].(map[string]any); got["running"] != true {
|
||||
t.Fatalf("an acknowledged call answered %v", got)
|
||||
}
|
||||
if c := l.Recent()[0]; c.State != CallFinishedAfter || !strings.Contains(string(c.Answer), "sent") {
|
||||
t.Fatalf("kept %+v", c)
|
||||
}
|
||||
}
|
||||
|
||||
// **A refused answer is recorded against its call, in the mesh's words** — not only as the client
|
||||
// library's line on standard error, which is all there was on 2026-10-05.
|
||||
func TestARefusedAnswerIsKeptAgainstItsCall(t *testing.T) {
|
||||
l, a := NewCallLog(), newAnswers(t)
|
||||
reply := "_INBOX.laptop.node-tools.jJ4DRJFnZYvkPUBKYwUG5v.LNRyztuX"
|
||||
l.serveCall("mesh-controller", "push", nil, reply, func(context.Context, json.RawMessage) (any, error) {
|
||||
return "told", nil
|
||||
}, a.respond, nil)
|
||||
var logged bytes.Buffer
|
||||
refusal := errors.New(`nats: permissions violation: Permissions Violation for Publish to "` + reply + `" on connection [838]`)
|
||||
if !l.Refusal(refusal, log.New(&logged, "", 0)) {
|
||||
t.Fatal("the refusal of a kept call's answer was not recognised")
|
||||
}
|
||||
c := l.Recent()[0]
|
||||
if c.Refused == "" || !strings.Contains(logged.String(), c.ID) {
|
||||
t.Fatalf("the refusal is not kept or not said: %+v / %q", c, logged.String())
|
||||
}
|
||||
other := errors.New(`nats: permissions violation: Permissions Violation for Publish to "mesh.node.x" on connection [1]`)
|
||||
if l.Refusal(other, nil) || l.Refusal(errors.New("nats: timeout"), nil) {
|
||||
t.Error("an unrelated error was taken for a refused answer")
|
||||
}
|
||||
}
|
||||
|
||||
// What `calls` keeps of the arguments never carries settings or a secret.
|
||||
func TestACallKeepsNoSettingsOrSecrets(t *testing.T) {
|
||||
got := string(kept(json.RawMessage(`{"module":"m","values":"{\"token\":\"s3cret\"}","secret":"s3cret"}`)))
|
||||
if strings.Contains(got, "s3cret") || !strings.Contains(got, `"module":"m"`) {
|
||||
t.Fatalf("kept %s", got)
|
||||
}
|
||||
}
|
||||
|
||||
// Only the newest KeptCalls are kept.
|
||||
func TestTheLogKeepsTheNewest(t *testing.T) {
|
||||
l := NewCallLog()
|
||||
for i := 0; i < KeptCalls+5; i++ {
|
||||
l.begin("s", "v", nil, "")
|
||||
}
|
||||
recent := l.Recent()
|
||||
if len(recent) != KeptCalls || !strings.HasSuffix(recent[0].ID, "-105") || !strings.HasSuffix(recent[KeptCalls-1].ID, "-6") {
|
||||
t.Fatalf("kept %d, newest %s, oldest %s", len(recent), recent[0].ID, recent[len(recent)-1].ID)
|
||||
}
|
||||
}
|
||||
@@ -25,8 +25,8 @@ type ToolHandler func(ctx context.Context, args json.RawMessage) (any, error)
|
||||
// SeatToolSubject is where a mesh-scoped seat's tool is asked (design 33 §4).
|
||||
func SeatToolSubject(seat, verb string) string { return "mesh.seat." + seat + ".tool." + verb }
|
||||
|
||||
// HandlerTimeout bounds one answer. A verb that runs a command — a push, a build with no wait —
|
||||
// answers in seconds; anything that has not in this long is said to have not answered.
|
||||
// HandlerTimeout bounds one call. Its caller is answered within AnswerWithin either way; this is how
|
||||
// long the call itself may run before it is stopped.
|
||||
const HandlerTimeout = 5 * time.Minute
|
||||
|
||||
// RebindAfter is how long a refused subscription waits before it is tried again.
|
||||
@@ -62,31 +62,17 @@ func (b OverNATS) serveTools(seat string, subjectOf func(string) string, handler
|
||||
_ = s.Unsubscribe()
|
||||
}
|
||||
}
|
||||
// A refused answer is recorded against its call, not only printed by the library.
|
||||
Calls.WatchRefusals(b.Conn, logger)
|
||||
for verb, handle := range handlers {
|
||||
verb, handle := verb, handle
|
||||
subject := subjectOf(verb)
|
||||
bind := func() (*nats.Subscription, error) {
|
||||
return b.Conn.QueueSubscribe(subject, "seat."+seat, func(msg *nats.Msg) {
|
||||
// Its own goroutine per call: a slow `push` must not hold up a `status` asked beside it,
|
||||
// and the library would otherwise run handlers one after another.
|
||||
go func() {
|
||||
ctx, cancel := context.WithTimeout(context.Background(), HandlerTimeout)
|
||||
defer cancel()
|
||||
args := json.RawMessage(msg.Data)
|
||||
if len(args) == 0 {
|
||||
args = json.RawMessage(`{}`)
|
||||
}
|
||||
var reply []byte
|
||||
result, err := handle(ctx, args)
|
||||
if err != nil {
|
||||
reply, _ = json.Marshal(map[string]any{"error": err.Error()})
|
||||
} else if reply, err = json.Marshal(map[string]any{"result": result}); err != nil {
|
||||
reply, _ = json.Marshal(map[string]any{"error": "the answer could not be written as JSON: " + err.Error()})
|
||||
}
|
||||
if err := msg.Respond(reply); err != nil && logger != nil {
|
||||
logger.Printf("%s: could not answer: %v", subject, err)
|
||||
}
|
||||
}()
|
||||
// and the library would otherwise run handlers one after another. Answered once, within
|
||||
// AnswerWithin, and kept (novox/hq issue 265).
|
||||
go Calls.serveCall(seat, verb, json.RawMessage(msg.Data), msg.Reply, handle, msg.Respond, logger)
|
||||
})
|
||||
}
|
||||
sub, err := bind()
|
||||
|
||||
Reference in New Issue
Block a user