package link import ( "context" "encoding/json" "errors" "fmt" "log" "regexp" "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 " 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 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"` // 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) } // 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..` // 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 } l.keep(*c) } // 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 } // 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 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 } // 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) // 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) 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) } }