From 801552c0eb5943ddf84f03a1d4541cc36a914c56 Mon Sep 17 00:00:00 2001 From: jochen Date: Tue, 6 Oct 2026 01:14:58 +0200 Subject: [PATCH] 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. --- cmd/mesh-controller/calls_verb_test.go | 81 +++++ cmd/mesh-controller/push.go | 2 + cmd/mesh-controller/seatverbs.go | 46 ++- cmd/mesh-controller/seatverbs_schema_test.go | 6 +- internal/broker/nats.go | 10 +- internal/catalogue/verbs.go | 11 +- internal/link/calls.go | 329 +++++++++++++++++++ internal/link/calls_test.go | 174 ++++++++++ internal/link/seattools.go | 28 +- 9 files changed, 658 insertions(+), 29 deletions(-) create mode 100644 cmd/mesh-controller/calls_verb_test.go create mode 100644 internal/link/calls.go create mode 100644 internal/link/calls_test.go diff --git a/cmd/mesh-controller/calls_verb_test.go b/cmd/mesh-controller/calls_verb_test.go new file mode 100644 index 0000000..56c3a19 --- /dev/null +++ b/cmd/mesh-controller/calls_verb_test.go @@ -0,0 +1,81 @@ +package main + +import ( + "context" + "encoding/json" + "strings" + "testing" + "time" + + "github.com/novox/mesh-controller/internal/broker" + "github.com/novox/mesh-controller/internal/link" +) + +// consoleWaits is how long the console waits for an answer (mesh-tools node-tools/internal/bus +// RequestTimeout) — the shortest wait of a caller the mesh ships. +const consoleWaits = 30 * time.Second + +// **A call's one answer is never later than its caller or the bus allow** (novox/hq issue 265): a +// holder answers within AnswerWithin, which must be inside both the console's wait and the window the +// bus gives an answer. Before, the console waited 30s, the bus 60s, and a push ran as long as it ran. +func TestAVerbAnswersInsideEveryWaitOnIt(t *testing.T) { + if link.AnswerWithin >= consoleWaits/2 { + t.Errorf("a call answers within %s: not well inside the console's %s", link.AnswerWithin, consoleWaits) + } + if link.AnswerWithin >= broker.ResponseTTL { + t.Errorf("a call answers within %s, after the bus stops permitting an answer at %s", link.AnswerWithin, broker.ResponseTTL) + } +} + +// A push — named or through command — answers before it runs: it sends the machine holding the bus +// first, and the broker reloading its user list forgets the answer it was about to permit. +func TestAPushAnswersBeforeItSends(t *testing.T) { + for _, c := range []struct { + verb string + args map[string]any + want bool + }{ + {"push", map[string]any{"node": "anchor"}, true}, + {"push", map[string]any{}, true}, + {"command", map[string]any{"command": "push anchor"}, true}, + {"command", map[string]any{"command": "push --behind"}, true}, + {"command", map[string]any{"command": "builds"}, false}, + {"status", map[string]any{}, false}, + {"assign", map[string]any{"node": "anchor", "module": "m"}, false}, + } { + argv, err := argvFor(c.verb, c.args) + if err != nil { + t.Fatalf("%s %v: %v", c.verb, c.args, err) + } + if got := answersFirst(argv); got != c.want { + t.Errorf("%s %v answers first: %v, want %v", c.verb, c.args, got, c.want) + } + } +} + +// `calls` is served, takes a call's id and nothing else, and says plainly when it holds no such call. +func TestCallsIsServedAndSaysWhatItKeeps(t *testing.T) { + handlers, behind, err := seatToolHandlers() + if err != nil || len(behind) != 0 { + t.Fatalf("%v %v", behind, err) + } + calls, ok := handlers["calls"] + if !ok { + t.Fatal("calls is not served") + } + if _, err := calls(context.Background(), json.RawMessage(`{"node":"anchor"}`)); err == nil || + !strings.Contains(err.Error(), `"node"`) { + t.Errorf("calls took an argument it does not declare: %v", err) + } + if _, err := calls(context.Background(), json.RawMessage(`{"call":"call-0-0"}`)); err == nil || + !strings.Contains(err.Error(), "not across a restart") { + t.Errorf("an unknown call was not said plainly: %v", err) + } + got, err := calls(context.Background(), json.RawMessage(`{}`)) + if err != nil { + t.Fatal(err) + } + if _, listed := got.(map[string]any)["calls"]; !listed { + t.Errorf("calls answered %v", got) + } +} diff --git a/cmd/mesh-controller/push.go b/cmd/mesh-controller/push.go index 011fbe7..b7fb7be 100644 --- a/cmd/mesh-controller/push.go +++ b/cmd/mesh-controller/push.go @@ -155,6 +155,8 @@ func serve(ctx context.Context) error { if !isNATS { return errors.New("the mesh's verbs are served over the bus, and this control plane is not on it") } + // A call that outlasts its caller's patience is followed by `calls` (novox/hq issue 265). + link.Calls.Follow = catalogue.ControllerSeatName + ".calls" stopServing, err := bus.ServeSeatTools(catalogue.ControllerSeatName, handlers, log.New(os.Stdout, "", log.LstdFlags)) if err != nil { return err diff --git a/cmd/mesh-controller/seatverbs.go b/cmd/mesh-controller/seatverbs.go index 245937e..25beacb 100644 --- a/cmd/mesh-controller/seatverbs.go +++ b/cmd/mesh-controller/seatverbs.go @@ -485,7 +485,7 @@ func seatToolHandlers() (map[string]link.ToolHandler, []string, error) { handlers := map[string]link.ToolHandler{} for _, v := range seat.Serves { verb := v.Name - if verb == "tools" { + if inProcess[verb] { handlers[verb] = func(ctx context.Context, raw json.RawMessage) (any, error) { args := map[string]any{} if len(bytes.TrimSpace(raw)) > 0 { @@ -493,10 +493,14 @@ func seatToolHandlers() (map[string]link.ToolHandler, []string, error) { return nil, fmt.Errorf("the arguments are not a JSON object: %w", err) } } - // Refused like any verb's: tools takes nothing, and something given is not ignored. - if _, err := readArguments(verb, args); err != nil { + // Refused like any verb's: what a verb does not take is not ignored. + a, err := readArguments(verb, args) + if err != nil { return nil, err } + if verb == "calls" { + return callsAnswer(link.Calls, a.given["call"]) + } return seatTools(), nil } continue @@ -535,12 +539,48 @@ func seatToolHandlers() (map[string]link.ToolHandler, []string, error) { if err != nil { return nil, err } + if answersFirst(argv) { + // Before anything is sent: a push sends the bus's own machine first, and a broker + // reloading its user list forgets the answer it was about to permit (novox/hq issue 265). + link.Acknowledge(ctx) + } return runVerb(ctx, argv) } } return handlers, behind, nil } +// inProcess are the verbs answered by this process rather than by a command it runs: `tools` from +// the records, `calls` from what this process served. +var inProcess = map[string]bool{"tools": true, "calls": true} + +// answersFirst is a command line whose caller is answered before it runs: a push, by its verb or +// through `command`. A push sends the machine holding the bus first when its user list changed, the +// broker reloads, and a reload forgets every answer the bus was about to permit — so an answer +// waiting for the push to end was refused, every time the list had changed (novox/hq issue 265). +func answersFirst(argv []string) bool { + return len(argv) > 0 && argv[0] == "push" +} + +// callsAnswer is what `calls` answers: the kept calls, newest first, without their answers — or +// one call whole. +func callsAnswer(log *link.CallLog, id string) (any, error) { + if id != "" { + c, ok := log.Get(id) + if !ok { + return nil, fmt.Errorf("no call %s is kept here: calls are kept by the controller that "+ + "answered them, the last %d, and not across a restart — `calls` lists them", id, link.KeptCalls) + } + return c, nil + } + recent := log.Recent() + for i := range recent { + recent[i].Answer = nil + } + return map[string]any{"calls": recent, "kept": link.KeptCalls, + "note": "newest first; `calls` with a call's id gives its whole answer"}, nil +} + // seatTools is what `tools` answers: every seat with a protocol, and the tools each serves, from the // mesh's own records — no holder in the path, so it is true while a holder restarts (design 33 §5). func seatTools() map[string]any { diff --git a/cmd/mesh-controller/seatverbs_schema_test.go b/cmd/mesh-controller/seatverbs_schema_test.go index b93a47b..d9a5d21 100644 --- a/cmd/mesh-controller/seatverbs_schema_test.go +++ b/cmd/mesh-controller/seatverbs_schema_test.go @@ -80,7 +80,7 @@ var saidOutright = map[string]string{ // shape of the push that named a machine and pushed every machine behind. func TestNoArgumentAVerbIsGivenIsPassedOver(t *testing.T) { for name, v := range servedSchemas(t) { - if name == "tools" { + if inProcess[name] { continue } names, switches := declaredArguments(v) @@ -123,7 +123,7 @@ func TestAnArgumentAVerbDoesNotDeclareIsRefused(t *testing.T) { args := sampleArguments(v) args[stranger] = "x" _, err := argvFor(name, args) - if name == "tools" { + if inProcess[name] { _, err = readArguments(name, args) } if err == nil || !strings.Contains(err.Error(), strconv.Quote(stranger)) { @@ -283,7 +283,7 @@ func TestEveryFlagOfAVerbsCommandIsAnArgumentOrAccountedFor(t *testing.T) { reached := map[string]map[string]bool{} // verb → flag sets its command lines reach served := servedSchemas(t) for name, v := range served { - if name == "tools" || name == "command" { + if inProcess[name] || name == "command" { continue } names, switches := declaredArguments(v) diff --git a/internal/broker/nats.go b/internal/broker/nats.go index 8c7d9b2..50ebab5 100644 --- a/internal/broker/nats.go +++ b/internal/broker/nats.go @@ -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") } diff --git a/internal/catalogue/verbs.go b/internal/catalogue/verbs.go index 98accf2..ae71f20 100644 --- a/internal/catalogue/verbs.go +++ b/internal/catalogue/verbs.go @@ -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", diff --git a/internal/link/calls.go b/internal/link/calls.go new file mode 100644 index 0000000..f094d26 --- /dev/null +++ b/internal/link/calls.go @@ -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 " 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) + } +} diff --git a/internal/link/calls_test.go b/internal/link/calls_test.go new file mode 100644 index 0000000..0d1f9b7 --- /dev/null +++ b/internal/link/calls_test.go @@ -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) + } +} diff --git a/internal/link/seattools.go b/internal/link/seattools.go index 89b24a5..10dc8dd 100644 --- a/internal/link/seattools.go +++ b/internal/link/seattools.go @@ -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()