diff --git a/cmd/mesh-controller/deadletters.go b/cmd/mesh-controller/deadletters.go index d2c7b735..c01dd2c4 100644 --- a/cmd/mesh-controller/deadletters.go +++ b/cmd/mesh-controller/deadletters.go @@ -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 diff --git a/cmd/mesh-controller/deadletters_test.go b/cmd/mesh-controller/deadletters_test.go index 5b3aef7a..e714456f 100644 --- a/cmd/mesh-controller/deadletters_test.go +++ b/cmd/mesh-controller/deadletters_test.go @@ -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 } diff --git a/cmd/mesh-controller/push.go b/cmd/mesh-controller/push.go index b4dabd7e..a4f54d63 100644 --- a/cmd/mesh-controller/push.go +++ b/cmd/mesh-controller/push.go @@ -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 diff --git a/cmd/mesh-controller/reconnects.go b/cmd/mesh-controller/reconnects.go index ed12cfe7..42f14479 100644 --- a/cmd/mesh-controller/reconnects.go +++ b/cmd/mesh-controller/reconnects.go @@ -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 -} diff --git a/cmd/mesh-controller/watchdogs.go b/cmd/mesh-controller/watchdogs.go index a4c07cef..88bf4434 100644 --- a/cmd/mesh-controller/watchdogs.go +++ b/cmd/mesh-controller/watchdogs.go @@ -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)