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) } }