From 5b7e6ff453eb991e21304d2fd31b4a13ae2c214b Mon Sep 17 00:00:00 2001 From: jochen Date: Thu, 8 Oct 2026 21:09:08 +0200 Subject: [PATCH] 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). --- cmd/mesh-controller/deadletters.go | 10 +++------- cmd/mesh-controller/deadletters_test.go | 6 +++--- cmd/mesh-controller/push.go | 2 ++ cmd/mesh-controller/reconnects.go | 14 +------------- cmd/mesh-controller/watchdogs.go | 4 ---- 5 files changed, 9 insertions(+), 27 deletions(-) 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)