Answer dead-letters on the lent serving connection, and keep one clip helper
Two handles to the same serving connection, and two copies of one helper, would drift (review of hq issues 327 and 330).
This commit is contained in:
@@ -6,7 +6,6 @@ import (
|
||||
"fmt"
|
||||
"strconv"
|
||||
"strings"
|
||||
"sync/atomic"
|
||||
|
||||
"github.com/nats-io/nats.go"
|
||||
|
||||
@@ -22,10 +21,6 @@ import (
|
||||
// serving controller is already on it. A verb run as a fresh process would open a connection of its own
|
||||
// for each call (novox/hq issue 327).
|
||||
|
||||
// deadLettersOn is the serving controller's bus, set when it starts watching the mesh; nil in a
|
||||
// process that serves nothing.
|
||||
var deadLettersOn atomic.Pointer[busHandles]
|
||||
|
||||
// busHandles are the serving controller's connection and JetStream handle.
|
||||
type busHandles struct {
|
||||
conn *nats.Conn
|
||||
@@ -41,11 +36,12 @@ const causeDeadLetter = "dead-letter"
|
||||
// deadLettersAnswer is what `dead-letters` answers: the list, one whole, or what came of delivering one
|
||||
// again or dropping it.
|
||||
func deadLettersAnswer(ctx context.Context, a *verbArguments) (any, error) {
|
||||
on := deadLettersOn.Load()
|
||||
if on == nil {
|
||||
serving := servingBus.Load()
|
||||
if serving == nil {
|
||||
return nil, errors.New("this controller is not serving, so it does not read DEAD_LETTERS: ask again, " +
|
||||
"and the serving controller answers")
|
||||
}
|
||||
on := &busHandles{conn: serving.Conn(), js: serving.Context()}
|
||||
// The shape first, and every argument it reads; one given beside it is refused before anything is
|
||||
// done, as every verb refuses what it would pass over (novox/hq issue 244).
|
||||
var deliver, drop, why, cause, idText, consumer, limit string
|
||||
|
||||
@@ -26,9 +26,9 @@ func servingDeadLetters(t *testing.T) *broker.JetStream {
|
||||
if err := js.EnsureControllerBuckets(); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
before := deadLettersOn.Load()
|
||||
deadLettersOn.Store(&busHandles{conn: js.Conn(), js: js.Context()})
|
||||
t.Cleanup(func() { deadLettersOn.Store(before) })
|
||||
before := servingBus.Load()
|
||||
servingBus.Store(js)
|
||||
t.Cleanup(func() { servingBus.Store(before) })
|
||||
return js
|
||||
}
|
||||
|
||||
|
||||
@@ -221,6 +221,8 @@ func serve(ctx context.Context) (err error) {
|
||||
handActConn = bus.Conn
|
||||
// And everything else this controller does on the bus for a moment (novox/hq issue 327).
|
||||
servingBus.Store(server.JetStream())
|
||||
// No longer serving: nothing is lent, and dead-letters says it is not read here (novox/hq issue 330).
|
||||
defer servingBus.Store(nil)
|
||||
// And says when it replaced a value given by hand (novox/hq ADR 0228).
|
||||
givenEvents = bus
|
||||
// And a pull request's merge check, asked when the forge announces its head and said when judged
|
||||
|
||||
@@ -143,7 +143,7 @@ func reconnecting(read closedConnections) []conditions.Observation {
|
||||
"client reconnecting in a loop", u.User, u.Dropped, partial, reconnectBound),
|
||||
Said: fmt.Sprintf("%d of %d closed connections dropped in %.0f h; by name %s; why %s", u.Dropped,
|
||||
u.Closed, hours, strings.Join(names, ", "), strings.Join(why, ", ")),
|
||||
Headline: clipWords(conditions.Capital(who)+" keeps losing the bus", 60),
|
||||
Headline: clip(conditions.Capital(who)+" keeps losing the bus", 60),
|
||||
Explanation: conditions.Capital(fmt.Sprintf("%s lost its connection to the bus %d times in the last hour "+
|
||||
"and connected again each time. While it reconnects, what it says and what it is asked waits.", who,
|
||||
u.Dropped)),
|
||||
@@ -171,15 +171,3 @@ func busUserWords(user string) (string, string) {
|
||||
}
|
||||
return "a client of the bus", ""
|
||||
}
|
||||
|
||||
// clipWords is words at most n characters long, cut at a word.
|
||||
func clipWords(s string, n int) string {
|
||||
if len(s) <= n {
|
||||
return s
|
||||
}
|
||||
cut := s[:n]
|
||||
if i := strings.LastIndex(cut, " "); i > 0 {
|
||||
cut = cut[:i]
|
||||
}
|
||||
return cut
|
||||
}
|
||||
|
||||
@@ -660,8 +660,6 @@ func watchTheMesh(ctx context.Context, open *stores, server *link.Server, bus li
|
||||
return nil
|
||||
}
|
||||
conditionsFrom = keeper
|
||||
// What a consumer gave up on is read and acted on by the verb in this process (novox/hq issue 330).
|
||||
deadLettersOn.Store(&busHandles{conn: server.JetStream().Conn(), js: server.JetStream().Context()})
|
||||
logf := func(format string, args ...any) { fmt.Printf(format+"\n", args...) }
|
||||
stopHearing, err := server.HearAdvisories(logf)
|
||||
if err != nil {
|
||||
@@ -688,8 +686,6 @@ func watchTheMesh(ctx context.Context, open *stores, server *link.Server, bus li
|
||||
return func() {
|
||||
stop()
|
||||
stopHearing()
|
||||
// No longer serving: the verb answers that it does not read DEAD_LETTERS here (novox/hq issue 330).
|
||||
deadLettersOn.Store(nil)
|
||||
flushing, cancel := context.WithTimeout(context.Background(), 10*time.Second)
|
||||
defer cancel()
|
||||
keeper.Close(flushing)
|
||||
|
||||
Reference in New Issue
Block a user