Files
mesh-controller/internal/link/calls.go
T
jochen 0c47f48537
mesh/merge-gate pass: builds build-agent, mesh-controller, route-proxy → ace, g14, novox, shanks; no bus step; every machine composes with the change as it…
mesh/repo-check pass: its merge-check.sh passed
mesh/delivery superseded: a newer head of the same pull request
mesh-cli lines are calls, followed to their answer; the terminal does not follow assign (review of hq ADR 0272)
- A line ran past the bus's one-minute window for an answer and its answer
  was refused; mesh-cli then said nothing ran of a line that had. Every
  line is now a call (calls.go): answered within AnswerWithin that it is
  running, with its call, and followed by the same account on the same
  node until it ends. Its record keeps the command word only, and its
  answer stays in the memory of the controller that ran it — never on the
  bus, never in `calls`, which answers anyone who may call the seat. Each
  line is said in the journal with its call, who asked where, and how.
- Lines running at once are bounded (8); one over is answered busy.
- An ordinary `settings` line is composed as the settings verb composes
  its own and meets that verb's refusals, the terminal-only keys among
  them, instead of the generic command's blanket refusal.
- The controller's module is assigned and unassigned at the terminal only:
  where it runs is the control-node, whose operator is the terminal.
- mesh.control.*.cli is in the writers table as the machine's own.
2026-10-09 12:48:12 +02:00

704 lines
26 KiB
Go

package link
import (
"context"
"encoding/json"
"errors"
"fmt"
"log"
"regexp"
"runtime/debug"
"sort"
"strconv"
"strings"
"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.
// ErrHandingOver is a call refused because the controller answering it is being replaced: stopping, or
// unable to run its own build because the node-engine has moved it aside (novox/hq issue 289). Nothing
// was done, and the controller after it answers the same call — so the answer carries
// `"retry": RetryHandingOver`, and a caller asks once more.
var ErrHandingOver = errors.New("the controller is handing over to the next one; nothing was done, ask again")
// RetryHandingOver is the `retry` mark of an answer refused by ErrHandingOver.
const RetryHandingOver = "handing-over"
// 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 this process keeps in memory, newest first — to match a refusal to its
// call, and to answer at once; the bus keeps the last thousand (Durably). keptAnswer is 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"
// CallAbandoned is a call whose controller stopped before it finished (novox/hq to-be 45 §6): a
// controller starting finds it running under another and says so, rather than leaving it running
// for ever in the record.
CallAbandoned = "abandoned: the controller running it stopped before it finished"
)
// 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"`
// Paged is an answer too large for one message, sent in parts to a caller that asked for them
// (novox/hq issue 314): how large, and until when it is held.
Paged string `json:"answer paged,omitempty"`
// Caller is the bus principal that asked, read from the inbox its answer went to — every principal
// is granted only its own (novox/hq to-be 45 §7: a hand act says who).
Caller string `json:"caller,omitempty"`
// Holder is the controller process that served it, so one starting can tell its own running
// calls from those a stopped one left.
Holder string `json:"holder,omitempty"`
// Epoch is the controller lease epoch it was served under (novox/hq to-be 45 §6); zero for none.
Epoch uint64 `json:"epoch,omitempty"`
reply string // the subject the answer went to, which the bus names when it refuses it
}
// CallKeeper keeps calls where the process serving them does not: the controller's bucket on the
// bus (novox/hq to-be 45 §6). A call is kept whole on every change — begun, finished, refused — so a
// controller replaced at any moment leaves the last word on each.
type CallKeeper interface {
Keep(ctx context.Context, c Call) error
// Kept is one call with its whole answer.
Kept(ctx context.Context, id string) (Call, bool, error)
// Recent is the kept calls, newest first, without their answers.
Recent(ctx context.Context) ([]Call, error)
}
// 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
// keeper keeps every call beyond this process, when Durably was given one; writes go through
// one goroutine, in order, so a call's last state is the one kept.
keeper CallKeeper
holder string
writes chan Call
logger *log.Logger
lost int // writes the keeper could not take, said once each
keepErr error
// epoch is the lease this process serves under (UnderLease): a call carries its epoch, and its
// record is written only while the lease is held. Nil writes every record with no epoch.
epoch func() (uint64, error)
// pages are the answers too large for one message, held for their callers (pages.go).
pages pageBook
}
// UnderLease makes every call carry the controller lease's epoch, and every write of a call's record
// pass the lease (novox/hq to-be 45 §6): a controller that lost the lease writes no finish over a call
// the next holder has marked abandoned — that mark is the record's last word.
func (l *CallLog) UnderLease(epoch func() (uint64, error)) {
l.mu.Lock()
defer l.mu.Unlock()
l.epoch = epoch
}
// 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} }
// keepTries is how many times one call's state is offered to the keeper before it is said lost: the
// bus reloading its user list refuses for a moment, and that is exactly when a push runs.
const keepTries = 5
// Durably keeps every call from now on with keeper as well as in memory, under this process's name,
// and marks running the calls a controller before this one left running: it stopped, so they cannot
// finish (novox/hq to-be 45 §6). Said, naming each.
func (l *CallLog) Durably(ctx context.Context, keeper CallKeeper, holder string, logger *log.Logger) error {
l.mu.Lock()
l.keeper, l.holder, l.logger = keeper, holder, logger
if l.writes == nil {
l.writes = make(chan Call, 256)
go l.keepWrites()
}
l.mu.Unlock()
kept, err := keeper.Recent(ctx)
if err != nil {
return fmt.Errorf("reading the calls kept on the bus: %w", err)
}
for _, c := range kept {
if c.State != CallRunning || c.Holder == holder {
continue
}
whole, found, err := keeper.Kept(ctx, c.ID)
if err != nil || !found {
whole = c
}
whole.State = CallAbandoned
if err := keeper.Keep(ctx, whole); err != nil {
return fmt.Errorf("marking %s abandoned: %w", c.ID, err)
}
if logger != nil {
logger.Printf("%s (%s.%s, asked %s by %s) was running under %s (epoch %d), which stopped: marked "+
"abandoned — it may have done part of what it was asked, and nothing will finish it", c.ID, c.Seat,
c.Verb, c.Started.Format(time.RFC3339), orSomebody(c.Caller), orSomebody(c.Holder), c.Epoch)
}
}
return nil
}
func orSomebody(s string) string {
if s == "" {
return "an unnamed caller"
}
return s
}
// keep queues one call's state for the keeper. Never blocks a call: a queue that is full is a keeper
// that is not taking writes, and that is said rather than waited on.
func (l *CallLog) keep(c Call) {
if l.writes == nil {
return
}
select {
case l.writes <- c:
default:
l.lost++
if l.logger != nil {
l.logger.Printf("%s (%s.%s) is kept in memory only: the bus is not taking calls' records (%d not kept)",
c.ID, c.Seat, c.Verb, l.lost)
}
}
}
func (l *CallLog) keepWrites() {
for c := range l.writes {
l.mu.Lock()
gate := l.epoch
l.mu.Unlock()
if gate != nil {
if _, err := gate(); err != nil {
if l.logger != nil {
l.logger.Printf("%s (%s.%s, %s) is not kept on the bus: %v — the controller holding the lease "+
"marks it abandoned", c.ID, c.Seat, c.Verb, c.State, err)
}
continue
}
}
var err error
for try := 0; try < keepTries; try++ {
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
err = l.keeper.Keep(ctx, c)
cancel()
if err == nil {
break
}
time.Sleep(time.Duration(try+1) * time.Second)
}
if err != nil && l.logger != nil {
l.logger.Printf("%s (%s.%s, %s) could not be kept on the bus after %d tries: %v — `calls` "+
"answers it from memory until this controller stops", c.ID, c.Seat, c.Verb, c.State, keepTries, err)
}
}
}
// callerOf is the bus principal an answer goes to: every principal's inbox is `_INBOX.<its user>.`
// followed by the client's own random token, and a user may itself hold dots.
func callerOf(reply string) string {
rest, ok := strings.CutPrefix(reply, "_INBOX.")
if !ok {
return ""
}
tokens := strings.Split(rest, ".")
for i, t := range tokens {
if i > 0 && isNUID(t) {
return strings.Join(tokens[:i], ".")
}
}
return ""
}
// isNUID is the client library's random inbox token: twenty-two letters and digits.
func isNUID(t string) bool {
if len(t) != 22 {
return false
}
for _, r := range t {
if !(r >= '0' && r <= '9' || r >= 'a' && r <= 'z' || r >= 'A' && r <= 'Z') {
return false
}
}
return true
}
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,
Caller: callerOf(reply), Holder: l.holder}
if l.epoch != nil {
c.Epoch, _ = l.epoch() // zero when not held: its record is then not written either
}
l.calls = append(l.calls, c)
if len(l.calls) > KeptCalls {
l.calls = l.calls[len(l.calls)-KeptCalls:]
}
l.keep(*c)
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
}
if shownOnce[c.Seat+"."+c.Verb] {
// **A secret shown once is never kept on the bus** (novox/hq ADR 0169): a join token is the
// one-time right to become a machine of the mesh, and `calls` answers anyone who may call the
// seat. Its caller was sent it; kept in this process's memory only when its caller has not
// had it yet — told the call was still running — and then until a restart, never beyond.
withheld, _ := json.Marshal(map[string]any{"not kept": "a secret shown once, to its caller"})
onTheBus := *c
onTheBus.Answer = withheld
if !answeredAlready {
c.Answer = withheld
}
l.keep(onTheBus)
return
}
onTheBus := *c
if c.Seat == CLISeat {
// **A mesh-cli line's answer is its asker's alone** (ADR 0272): what a command at the controller's terminal
// printed — a token, a secret's reference — and `calls` answers anyone who may call the seat. Kept in this
// process's memory for the asker to follow, and never on the bus.
onTheBus.Answer, _ = json.Marshal(map[string]any{"not kept": "a mesh-cli line's answer, its asker's alone: " +
"kept in the memory of the controller that ran it, for the asker to follow"})
l.keep(onTheBus)
return
}
if len(onTheBus.Answer) > keptOnTheBusAtMost {
// A record on the bus is one message too, and one larger than the bus carries is refused whole —
// the call's state with it (novox/hq issue 314). The answer stays in this process's memory, and
// `calls` with its id pages it from there.
onTheBus.Answer, _ = json.Marshal(map[string]any{"not kept": fmt.Sprintf("an answer of %d bytes, more than "+
"a record on the bus holds: kept in the memory of the controller that answered it", len(c.Answer))})
}
l.keep(onTheBus)
}
// keptOnTheBusAtMost is the largest answer a call's record on the bus carries: half of the bus's default
// message limit, leaving the rest of the record room.
const keptOnTheBusAtMost = 512 << 10
// shownOnce are the verbs whose answer is a secret shown once to its caller, as `<seat>.<verb>`.
var shownOnce = map[string]bool{"mesh-controller.token": true}
// Recent is the kept calls, newest first, as copies: this process's from memory, and, when they are
// kept durably, every other the bus holds — a call a controller before this one served included.
// Memory wins for a call in both, being the newer word on it. A bus that cannot be read is said in
// the error beside what memory holds, never answered as no calls.
// Running is every call this process is serving that has not finished, oldest first, without
// answers: what the watchdog of a call's bound (novox/hq to-be 45 S7) reads. This process's own,
// because a call another controller left running is said abandoned when this one starts.
func (l *CallLog) Running() []Call {
l.mu.Lock()
defer l.mu.Unlock()
var out []Call
for _, c := range l.calls {
if c.State == CallRunning {
running := *c
running.Answer = nil
out = append(out, running)
}
}
return out
}
func (l *CallLog) Recent() ([]Call, error) {
l.mu.Lock()
out := make([]Call, 0, len(l.calls))
seen := map[string]bool{}
for i := len(l.calls) - 1; i >= 0; i-- {
out = append(out, *l.calls[i])
seen[l.calls[i].ID] = true
}
keeper := l.keeper
l.mu.Unlock()
if keeper == nil {
return out, nil
}
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
kept, err := keeper.Recent(ctx)
if err != nil {
return out, fmt.Errorf("the calls kept on the bus could not be read, so only this controller's own "+
"are listed: %w", err)
}
for _, c := range kept {
if !seen[c.ID] {
out = append(out, c)
}
}
sort.SliceStable(out, func(i, j int) bool { return out[i].Started.After(out[j].Started) })
return out, nil
}
// inMemory is one call this process served and still holds, whole.
func (l *CallLog) inMemory(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
}
// Shown is one kept call as `calls` shows it: whole, but for a mesh-cli line's answer, which is its asker's alone
// and followed by it (ADR 0272) — `calls` answers anyone who may call the seat.
func (l *CallLog) Shown(id string) (Call, bool, error) {
c, found, err := l.Get(id)
if found && c.Seat == CLISeat {
c.Answer, _ = json.Marshal(map[string]any{"not shown": "a mesh-cli line's answer is its asker's alone: " +
"mesh-cli follows it"})
}
return c, found, err
}
// Get is one kept call: from memory, or from the bus when this process did not serve it.
func (l *CallLog) Get(id string) (Call, bool, error) {
l.mu.Lock()
for _, c := range l.calls {
if c.ID == id {
found := *c
l.mu.Unlock()
return found, true, nil
}
}
keeper := l.keeper
l.mu.Unlock()
if keeper == nil {
return Call{}, false, nil
}
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
return keeper.Kept(ctx, id)
}
// IsDurable says whether calls outlive this process.
func (l *CallLog) IsDurable() bool {
l.mu.Lock()
defer l.mu.Unlock()
return l.keeper != nil
}
// 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 k == "line":
// A mesh-cli line (ADR 0272): its command word, never the rest, which may carry settings.
if words, ok := v.([]any); ok && len(words) > 0 {
out[k] = []any{words[0], "(the rest 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.keep(*hit)
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
}
type callIDKey struct{}
// WithCallID is ctx carrying the call it serves; CallIDIn reads it back, empty outside one.
func WithCallID(ctx context.Context, id string) context.Context {
return context.WithValue(ctx, callIDKey{}, id)
}
// CallIDIn is the call ctx serves, or empty.
func CallIDIn(ctx context.Context) string {
id, _ := ctx.Value(callIDKey{}).(string)
return id
}
// 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) {
l.serveCallWithin(seat, verb, args, reply, handle, respond, 0, logger)
}
// serveCallWithin is serveCall on a bus that carries at most limit bytes in one message: an answer
// larger than that is paged (pages.go), and one the bus would not carry anyway is said to its caller
// and to `calls` — never dropped with only a line in the journal (novox/hq issue 314). A limit of 0 is
// one not known.
func (l *CallLog) serveCallWithin(seat, verb string, args json.RawMessage, reply string, handle ToolHandler,
respond func([]byte) error, limit int64, 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)
// Who asked travels with the call, so an act it does by hand says so (novox/hq to-be 45 §7).
ctx = context.WithValue(ctx, callerKey{}, c.Caller)
// And which call it is, for what the handler says of it in the journal.
ctx = WithCallID(ctx, c.ID)
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
// **A handler that panics is an answer, not a dead controller.** A verb answered in the serving
// process (novox/hq issue 327: conditions, hand-acts, queue, dead-letters) runs on this goroutine,
// where a panic would take every other call and the controller with it.
defer func() {
if p := recover(); p != nil {
if logger != nil {
logger.Printf("%s.%s: call %s panicked: %v\n%s", seat, verb, c.ID, p, debug.Stack())
}
body, _ = json.Marshal(map[string]any{"error": fmt.Sprintf("%s failed inside the controller and "+
"answered nothing: %v. It is in the controller's log", verb, p)})
done <- outcome{body, true}
}
}()
result, err := handle(ctx, args)
failed := err != nil
if errors.Is(err, ErrHandingOver) {
// Marked as well as said, so a caller asks again without reading the words (issue 289).
body, _ = json.Marshal(map[string]any{"error": err.Error(), "retry": RetryHandingOver})
} else 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 limit > 0 && int64(len(body)) > limit {
if logger != nil {
logger.Printf("%s.%s: the answer to %s is %d bytes, more than the bus carries in one message (%d): "+
"paged", seat, verb, c.ID, len(body), limit)
}
body = l.page(c, body, limit)
}
err := respond(body)
if err == nil {
return
}
if logger != nil {
logger.Printf("%s.%s: could not answer %s: %v", seat, verb, c.ID, err)
}
l.unsent(c, err)
if errors.Is(err, nats.ErrMaxPayload) {
// The limit was not known, or was not the bus's: its caller is told, in a message that fits.
if again := respond(cannotCarry(seat, verb, c.ID, len(body), err)); again != nil && logger != nil {
logger.Printf("%s.%s: could not say to the caller of %s why it has no answer: %v", seat, verb, c.ID, again)
}
}
}
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)
}
}