diff --git a/cmd/mesh-controller/deadletters.go b/cmd/mesh-controller/deadletters.go index d929bb9a..d2c7b735 100644 --- a/cmd/mesh-controller/deadletters.go +++ b/cmd/mesh-controller/deadletters.go @@ -125,7 +125,7 @@ func actOnDeadLetter(ctx context.Context, on *busHandles, act string, id uint64, } if act == "deliver" { // Refused before it is recorded: an act that cannot be done is not an act. - if _, err := link.AgainTo(d); err != nil { + if _, err := link.AgainTo(on.js, d); err != nil { return nil, err } } diff --git a/cmd/mesh-controller/probes.go b/cmd/mesh-controller/probes.go index 1f19b088..6e034d07 100644 --- a/cmd/mesh-controller/probes.go +++ b/cmd/mesh-controller/probes.go @@ -878,6 +878,10 @@ func streamDiffers(want broker.Stream, have nats.StreamConfig) string { if want.MaxBytes > 0 && have.MaxBytes != want.MaxBytes { differs = append(differs, fmt.Sprintf("holds up to %d bytes, defined %d", have.MaxBytes, want.MaxBytes)) } + if want.DuplicatesSeconds > 0 && have.Duplicates != time.Duration(want.DuplicatesSeconds)*time.Second { + differs = append(differs, fmt.Sprintf("keeps one of a message id for %s, defined %s", have.Duplicates, + time.Duration(want.DuplicatesSeconds)*time.Second)) + } if want.DiscardNew && have.Discard != nats.DiscardNew { differs = append(differs, "drops what it holds when full, defined to refuse what comes next") } diff --git a/cmd/mesh-controller/watchdogs.go b/cmd/mesh-controller/watchdogs.go index 1bede623..a4c07cef 100644 --- a/cmd/mesh-controller/watchdogs.go +++ b/cmd/mesh-controller/watchdogs.go @@ -688,6 +688,8 @@ 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) diff --git a/internal/broker/nats.go b/internal/broker/nats.go index d313237a..3e765d71 100644 --- a/internal/broker/nats.go +++ b/internal/broker/nats.go @@ -413,11 +413,10 @@ func PermissionsFor(p Principal) (Permissions, error) { // Named, not `$JS.EVENT.>`: the other advisories are every API call the mesh makes. sub = append(sub, BusAdvisories...) // **And what a consumer gave up on, kept and delivered again** (novox/hq issue 330): the notice - // acknowledged once the message is copied, the copy kept, and a message delivered again — an - // event to the one consumer that gave it up, an ask back onto its seat's queue, whose one worker - // is the consumer that gave it up. + // acknowledged once the message is copied, the copy kept, and an event delivered again to the one + // consumer that gave it up. An ask to a seat is not delivered again, so no seat's queue is granted. pub = append(pub, "$JS.ACK."+DeadLetterNoticesStream+"."+ControllerName+".>", deadLetterPrefix+">", - againPrefix+">", "mesh.seat.*.accept.>") + againPrefix+">") case KindPerson: // Tools, and nothing else. Every subject a person may publish is a tool call; a person diff --git a/internal/broker/streams.go b/internal/broker/streams.go index 2e514e47..f1bc66bf 100644 --- a/internal/broker/streams.go +++ b/internal/broker/streams.go @@ -440,6 +440,31 @@ func MeshConsumers() []Consumer { } } +// NoticesConsumer is the controller's consumer on DEAD_LETTER_NOTICES (novox/hq issue 330). +func NoticesConsumer() Consumer { + for _, c := range MeshConsumers() { + if c.Stream == DeadLetterNoticesStream { + return c + } + } + panic("the mesh's consumers carry none on " + DeadLetterNoticesStream) +} + +// AssertServingConsumers are the controller's own consumers its serving cannot go without: all but the +// one on DEAD_LETTER_NOTICES, which the keeper of dead letters asserts and retries by itself, so a fault +// there never stops the controller serving (novox/hq issue 330). +func AssertServingConsumers(e Ensurer) error { + for _, c := range MeshConsumers() { + if c.Stream == DeadLetterNoticesStream { + continue + } + if err := e.EnsureConsumer(c); err != nil { + return fmt.Errorf("asserting consumer %s on %s: %w", c.Name, c.Stream, err) + } + } + return nil +} + // Ensurer is the part of a JetStream connection consumer assertion needs, narrow for the reason // Asserter is. type Ensurer interface { diff --git a/internal/broker/testdata/composed.conf b/internal/broker/testdata/composed.conf index 5f31c221..0519c572 100644 --- a/internal/broker/testdata/composed.conf +++ b/internal/broker/testdata/composed.conf @@ -24,7 +24,7 @@ accounts { jetstream: enabled users = [ { user: "controller", password: "$2a$11$cccccccccccccccccccccc", permissions: { - publish: { allow: ["$JS.ACK.CONTROL.controller.>", "$JS.ACK.DEAD_LETTER_NOTICES.controller.>", "$JS.ACK.EVENTS.controller.>", "$JS.API.>", "$KV.SEAT_MESH_BUILD_MACHINE_cancelled.>", "$KV.SEAT_NODE_BUILD_AGENT_cancelled.>", "$KV.mesh-controller_calls.>", "$KV.mesh-controller_condition-history.>", "$KV.mesh-controller_conditions.>", "$KV.mesh-controller_hand-acts.>", "$KV.mesh-controller_lease.>", "$SRV.INFO", "_INBOX.enrol.>", "mesh.again.>", "mesh.assignment.>", "mesh.events.dead.>", "mesh.mod.*.tool.>", "mesh.node.>", "mesh.seat.*.accept.>", "mesh.seat.mesh-build-machine.accept.>", "mesh.seat.mesh-build-machine.tool.>", "mesh.seat.mesh-controller.event.applied", "mesh.seat.mesh-controller.event.built-before", "mesh.seat.mesh-controller.event.checked", "mesh.seat.mesh-controller.event.condition-changed", "mesh.seat.mesh-controller.event.condition-cleared", "mesh.seat.mesh-controller.event.condition-raised", "mesh.seat.mesh-controller.event.doctor-heartbeat", "mesh.seat.mesh-controller.event.healer-acted", "mesh.seat.mesh-controller.event.plan-moved", "mesh.seat.mesh-controller.event.refused", "mesh.seat.mesh-controller.event.rolled-back", "mesh.seat.mesh-controller.event.secret-replaced", "mesh.seat.mesh-delivery.tool.close", "mesh.seat.mesh-delivery.tool.stalled", "mesh.seat.node-backup.tool.backed-up.*", "mesh.seat.node-backup.tool.now.*", "mesh.seat.node-build-agent.accept.>", "mesh.seat.node-build-agent.tool.>", "mesh.seat.node-intrusion-prevention.tool.banned.*"] } + publish: { allow: ["$JS.ACK.CONTROL.controller.>", "$JS.ACK.DEAD_LETTER_NOTICES.controller.>", "$JS.ACK.EVENTS.controller.>", "$JS.API.>", "$KV.SEAT_MESH_BUILD_MACHINE_cancelled.>", "$KV.SEAT_NODE_BUILD_AGENT_cancelled.>", "$KV.mesh-controller_calls.>", "$KV.mesh-controller_condition-history.>", "$KV.mesh-controller_conditions.>", "$KV.mesh-controller_hand-acts.>", "$KV.mesh-controller_lease.>", "$SRV.INFO", "_INBOX.enrol.>", "mesh.again.>", "mesh.assignment.>", "mesh.events.dead.>", "mesh.mod.*.tool.>", "mesh.node.>", "mesh.seat.mesh-build-machine.accept.>", "mesh.seat.mesh-build-machine.tool.>", "mesh.seat.mesh-controller.event.applied", "mesh.seat.mesh-controller.event.built-before", "mesh.seat.mesh-controller.event.checked", "mesh.seat.mesh-controller.event.condition-changed", "mesh.seat.mesh-controller.event.condition-cleared", "mesh.seat.mesh-controller.event.condition-raised", "mesh.seat.mesh-controller.event.doctor-heartbeat", "mesh.seat.mesh-controller.event.healer-acted", "mesh.seat.mesh-controller.event.plan-moved", "mesh.seat.mesh-controller.event.refused", "mesh.seat.mesh-controller.event.rolled-back", "mesh.seat.mesh-controller.event.secret-replaced", "mesh.seat.mesh-delivery.tool.close", "mesh.seat.mesh-delivery.tool.stalled", "mesh.seat.node-backup.tool.backed-up.*", "mesh.seat.node-backup.tool.now.*", "mesh.seat.node-build-agent.accept.>", "mesh.seat.node-build-agent.tool.>", "mesh.seat.node-intrusion-prevention.tool.banned.*"] } subscribe: { allow: ["$JS.API.>", "$JS.EVENT.ADVISORY.CONSUMER.DELETED.>", "$JS.EVENT.ADVISORY.CONSUMER.MAX_DELIVERIES.>", "$SRV.INFO", "$SRV.INFO.mesh-controller", "$SRV.INFO.mesh-controller.>", "$SRV.PING", "$SRV.PING.mesh-controller", "$SRV.PING.mesh-controller.>", "$SRV.STATS", "$SRV.STATS.mesh-controller", "$SRV.STATS.mesh-controller.>", "_DELIVER.controller", "_DELIVER.controller.>", "_INBOX.controller.>", "mesh.control.>", "mesh.mod.*.event.provisioner.failing", "mesh.mod.*.event.provisioner.recovered", "mesh.mod.*.event.provisioner.retirement", "mesh.mod.gitea.event.pull.merged", "mesh.mod.gitea.event.pull.updated", "mesh.mod.mesh-catalog.event.catching-up", "mesh.mod.mesh-catalog.event.upgraded", "mesh.seat.mesh-build-machine.event.built", "mesh.seat.mesh-controller.tool.>", "mesh.seat.node-build-agent.event.built"] } allow_responses: { max: 1, ttl: "1m" } } } diff --git a/internal/link/bus_nats.go b/internal/link/bus_nats.go index b902da42..bfaebdb4 100644 --- a/internal/link/bus_nats.go +++ b/internal/link/bus_nats.go @@ -8,13 +8,16 @@ import ( // JetStream connection the caller has already raised the streams on. Nothing is declared here — // the streams and the controller's consumers are asserted by Raise, before anything is served. func ConnectNats(js *broker.JetStream, enroller Enroller, listener Listener) *Server { + logger := newLog() + inbound := Nats(js).(*natsInbound) + inbound.log = logger return &Server{ - inbound: Nats(js), + inbound: inbound, bus: OverNATS{JS: js.Context(), Conn: js.Conn()}, js: js, enroller: enroller, listener: listener, - log: newLog(), + log: logger, } } diff --git a/internal/link/deadletters.go b/internal/link/deadletters.go index 6cc6fd10..8c0a4eaf 100644 --- a/internal/link/deadletters.go +++ b/internal/link/deadletters.go @@ -1,10 +1,12 @@ package link import ( + "context" "encoding/json" "errors" "fmt" "log" + "slices" "strconv" "strings" "time" @@ -60,9 +62,10 @@ type DeadLetter struct { // taken. The record of it is kept all the same, so it is said and dropped, never silently missing. Lost string `json:"lost,omitempty"` Size int `json:"size"` - // Body and Headers are the message itself, given only for one dead letter asked by its id. - Body string `json:"body,omitempty"` - Headers map[string]string `json:"headers,omitempty"` + // Body and Headers are the message itself, given only for one dead letter asked by its id; a header + // with several values keeps them all. + Body string `json:"body,omitempty"` + Headers map[string][]string `json:"headers,omitempty"` } // maxDeliveries is the part of the server's notice a kept message is made from. @@ -139,10 +142,10 @@ func deadLetterOf(id uint64, subject string, h nats.Header, data []byte, whole b d.Published, _ = time.Parse(time.RFC3339Nano, h.Get(DeadHeaderPublished)) if whole { d.Body = string(data) - d.Headers = map[string]string{} - for k := range h { + d.Headers = map[string][]string{} + for k, v := range h { if !strings.HasPrefix(k, "Mesh-Dead-") && !strings.HasPrefix(k, "Nats-") { - d.Headers[k] = h.Get(k) + d.Headers[k] = append([]string(nil), v...) } } } @@ -164,16 +167,32 @@ func HeldDeadLetters(js nats.JetStreamContext) (map[string]int, error) { return held, nil } -// DeadLetters lists what DEAD_LETTERS holds, newest first, at most most of them; for one consumer when -// consumer names one (its name, or `.`). +// DeadLetters lists what DEAD_LETTERS holds, newest first, at most most of them (all when most is not +// positive); for one consumer when consumer names one (its name, or `.`). The total is +// the stream's own count per consumer, so it is right however few are read. func DeadLetters(js nats.JetStreamContext, consumer string, most int) ([]DeadLetter, int, error) { + held, err := HeldDeadLetters(js) + if err != nil { + return nil, 0, err + } + total := 0 + for key, n := range held { + if consumer == "" || key == consumer || strings.HasSuffix(key, "."+consumer) { + total += n + } + } + if total == 0 { + return nil, 0, nil + } info, err := js.StreamInfo(broker.DeadLettersStream) if err != nil { return nil, 0, fmt.Errorf("%s cannot be read: %w", broker.DeadLettersStream, err) } var out []DeadLetter - total := 0 for seq := info.State.LastSeq; seq >= info.State.FirstSeq && seq > 0; seq-- { + if most > 0 && len(out) >= most || len(out) >= total { + break + } raw, err := js.GetMsg(broker.DeadLettersStream, seq) if errors.Is(err, nats.ErrMsgNotFound) { continue // delivered again or dropped @@ -185,10 +204,7 @@ func DeadLetters(js nats.JetStreamContext, consumer string, most int) ([]DeadLet if consumer != "" && consumer != d.Consumer && consumer != d.Stream+"."+d.Consumer { continue } - total++ - if most <= 0 || len(out) < most { - out = append(out, d) - } + out = append(out, d) } return out, total, nil } @@ -210,23 +226,38 @@ func DeadLetterNamed(js nats.JetStreamContext, id uint64) (DeadLetter, error) { } // AgainTo is where a kept message is delivered again so that only the consumer that gave it up gets -// it: an event under that consumer's own again subject on EVENTS; an ask back onto its seat's queue, -// whose one worker is that consumer. Any other stream's message is refused, with why. -func AgainTo(d DeadLetter) (string, error) { +// it: an event under that consumer's own again subject on EVENTS — and only when the consumer exists and +// filters that subject, so a message is never let go as delivered while nobody receives it. Any other +// stream's message is refused, with why: an ask given up on by a seat's worker is not delivered again yet +// (novox/hq issue 330's follow-up), since publishing it again leaves the original stuck in the queue. +func AgainTo(js nats.JetStreamContext, d DeadLetter) (string, error) { switch { case d.Lost != "": return "", fmt.Errorf("dead letter %d holds no message to deliver: %s. Drop it", d.ID, d.Lost) case d.Subject == "": return "", fmt.Errorf("dead letter %d does not say the subject it was published on, so it cannot be "+ "delivered again. Drop it", d.ID) - case d.Stream == broker.EventsStream: - return broker.AgainSubject(d.Consumer, d.Subject), nil - case strings.HasPrefix(d.Stream, "SEAT_") && strings.HasSuffix(d.Consumer, "_worker"): - return d.Subject, nil + case d.Stream != broker.EventsStream: + return "", fmt.Errorf("dead letter %d is from %s, and only an event is delivered again: publishing it "+ + "again would reach every consumer of its subject, or leave the original in its queue. Drop it, and "+ + "have its sender say it again", d.ID, d.Stream) } - return "", fmt.Errorf("dead letter %d is from %s, which nothing delivers again to one consumer alone: "+ - "publishing it again would reach every consumer of its subject. Drop it, and have its sender say it again", - d.ID, d.Stream) + info, err := js.ConsumerInfo(d.Stream, d.Consumer) + if errors.Is(err, nats.ErrConsumerNotFound) { + return "", fmt.Errorf("dead letter %d was given up by %s, which is no longer on the bus, so nothing would "+ + "receive it. Drop it, or deliver it again once the module is assigned there again", d.ID, d.Who) + } + if err != nil { + return "", fmt.Errorf("whether %s can receive dead letter %d cannot be read: %w", d.Who, d.ID, err) + } + want := broker.AgainFilter(d.Consumer) + filters := append([]string{info.Config.FilterSubject}, info.Config.FilterSubjects...) + if !slices.Contains(filters, want) { + return "", fmt.Errorf("%s does not yet take events delivered again (it does not filter %s): the controller "+ + "sets that at the next send to its machine. Nothing was done, and dead letter %d is still kept", + d.Who, want, d.ID) + } + return broker.AgainSubject(d.Consumer, d.Subject), nil } // DeliverAgain hands a kept message to the consumer that gave it up, and nobody else, then removes it @@ -237,13 +268,13 @@ func DeliverAgain(js nats.JetStreamContext, id uint64) (DeadLetter, string, erro if err != nil { return d, "", err } - to, err := AgainTo(d) + to, err := AgainTo(js, d) if err != nil { return d, "", err } again := &nats.Msg{Subject: to, Header: nats.Header{}, Data: []byte(d.Body)} for k, v := range d.Headers { - again.Header.Set(k, v) + again.Header[k] = append([]string(nil), v...) } again.Header.Set(AgainHeader, strconv.FormatUint(id, 10)) again.Header.Set(nats.MsgIdHdr, "again."+broker.DeadLettersStream+"."+strconv.FormatUint(id, 10)) @@ -272,6 +303,45 @@ func DropDeadLetter(js nats.JetStreamContext, id uint64) (DeadLetter, error) { // noticeRetry is how long a notice whose message could not be kept waits before it is taken again. var noticeRetry = time.Minute +// keepingGivenUp keeps taking the notices until ctx ends; the returned function stops it. A subscription +// that cannot be made is said — in the log, and as the max-deliveries condition of the notices themselves +// — and tried again every noticeRetry, so it never stops the controller serving. +func keepingGivenUp(ctx context.Context, bus *broker.JetStream, logger *log.Logger) func() { + ctx, cancel := context.WithCancel(ctx) + done := make(chan struct{}) + go func() { + defer close(done) + for { + sub, err := func() (*nats.Subscription, error) { + if err := bus.EnsureConsumer(broker.NoticesConsumer()); err != nil { + return nil, err + } + return keepGivenUp(bus.Context(), logger) + }() + if err == nil { + <-ctx.Done() + _ = sub.Unsubscribe() + return + } + logger.Printf("the notices of messages consumers gave up on cannot be taken, so none is kept until "+ + "they can: %v", err) + Advisories.Heard(Advisory{Kind: AdvisoryMaxDeliveries, ID: broker.DeadLetterNoticesStream + "." + + broker.ControllerName, Stream: broker.DeadLetterNoticesStream, Consumer: broker.ControllerName, + Token: AdvisoryNotKept, Said: "the controller cannot take the notices of messages consumers gave " + + "up on, so none is kept: " + err.Error()}, time.Now()) + select { + case <-ctx.Done(): + return + case <-time.After(noticeRetry): + } + } + }() + return func() { + cancel() + <-done + } +} + // keepGivenUp takes the server's maximum-deliveries notices off their stream and keeps the message // each is about, for as long as the subscription stands. A notice that could not be kept is offered // again after noticeRetry, and said: in the log, and as the consumer's max-deliveries condition. diff --git a/internal/link/deadletters_test.go b/internal/link/deadletters_test.go index ea013d09..12928254 100644 --- a/internal/link/deadletters_test.go +++ b/internal/link/deadletters_test.go @@ -4,6 +4,7 @@ import ( "context" "encoding/json" "errors" + "fmt" "io" "log" "strings" @@ -113,7 +114,7 @@ func TestADeadLetterIsDeliveredAgainToItsConsumerAlone(t *testing.T) { } whole, err := DeadLetterNamed(js.Context(), d.ID) - if err != nil || whole.Body != `{"n":1}` || whole.Headers["x-node"] != "forge" { + if err != nil || whole.Body != `{"n":1}` || len(whole.Headers["x-node"]) != 1 || whole.Headers["x-node"][0] != "forge" { t.Fatalf("one asked whole is %+v (%v)", whole, err) } @@ -218,22 +219,68 @@ func TestASubjectDeliveredAgainSaysTheEventsOwn(t *testing.T) { } } -// A message from a stream no consumer alone can be given again on is refused, with why. -func TestOnlyAnEventOrAnAskIsDeliveredAgain(t *testing.T) { - for _, c := range []struct { - d DeadLetter - want string - }{ - {DeadLetter{ID: 1, Stream: "EVENTS", Consumer: "a_b", Subject: "mesh.mod.x.event.y"}, "mesh.again.a_b.mod.x.event.y"}, - {DeadLetter{ID: 2, Stream: "SEAT_TELEGRAM_SENDER", Consumer: "SEAT_TELEGRAM_SENDER_worker", - Subject: "mesh.seat.telegram-sender.accept.send"}, "mesh.seat.telegram-sender.accept.send"}, +// Only an event is delivered again, and only to a consumer that is on the bus and filters its again +// subject: a dead letter is never let go as delivered while nobody receives it. +func TestADeadLetterIsDeliveredAgainOnlyWhereItIsReceived(t *testing.T) { + js := aBus(t) + keeping(t, js) + failing := aModuleConsumer(t, js, "media", "sonarr", "gitea.pull.merged") + if _, err := js.Context().Publish("mesh.mod.gitea.event.pull.merged", []byte(`{}`)); err != nil { + t.Fatal(err) + } + d := givenUp(t, js, failing) + + // A consumer as it was before this fix: its filters do not take what is delivered again. + c, _ := broker.ConsumerFor(broker.Principal{Kind: broker.KindModule, Node: "media", Module: "sonarr", + Consumes: []string{"gitea.pull.merged"}, PasswordHash: "x"}) + c.Filters = c.Filters[:len(c.Filters)-1] + if err := js.EnsureConsumer(c); err != nil { + t.Fatal(err) + } + if _, _, err := DeliverAgain(js.Context(), d.ID); err == nil || !strings.Contains(err.Error(), "does not yet take") { + t.Fatalf("delivered to a consumer that does not filter its again subject: %v", err) + } + if _, err := DeadLetterNamed(js.Context(), d.ID); err != nil { + t.Fatalf("a refused delivery let the dead letter go: %v", err) + } + + // The consumer gone: nothing would receive it. + if err := js.Context().DeleteConsumer(broker.EventsStream, c.Name); err != nil { + t.Fatal(err) + } + if _, _, err := DeliverAgain(js.Context(), d.ID); err == nil || !strings.Contains(err.Error(), "no longer on the bus") { + t.Fatalf("delivered to a consumer that is gone: %v", err) + } + if _, err := DeadLetterNamed(js.Context(), d.ID); err != nil { + t.Fatalf("a refused delivery let the dead letter go: %v", err) + } + + // Not an event: refused, whatever stream it is from. + for _, other := range []DeadLetter{ + {ID: 2, Stream: "SEAT_TELEGRAM_SENDER", Consumer: "SEAT_TELEGRAM_SENDER_worker", Subject: "mesh.seat.telegram-sender.accept.send"}, + {ID: 3, Stream: "KV_x", Consumer: "y", Subject: "$KV.x.k"}, } { - if to, err := AgainTo(c.d); err != nil || to != c.want { - t.Errorf("%+v: %s %v", c.d, to, err) + if _, err := AgainTo(js.Context(), other); err == nil { + t.Errorf("%s was given a subject to be delivered again on", other.Stream) } } - if _, err := AgainTo(DeadLetter{ID: 3, Stream: "KV_x", Consumer: "y", Subject: "$KV.x.k"}); err == nil { - t.Error("a bucket's message was given a subject to be delivered again on") +} + +// A total says how many are held, however few the list carries. +func TestTheListSaysTheTotalAndStopsAtItsLimit(t *testing.T) { + js := aBus(t) + for seq := 1; seq <= 5; seq++ { + if _, err := KeepDeadLetter(js.Context(), []byte(fmt.Sprintf(`{"stream":"EVENTS","consumer":"media_sonarr",`+ + `"stream_seq":%d,"deliveries":5}`, seq))); err != nil { + t.Fatal(err) + } + } + list, total, err := DeadLetters(js.Context(), "media_sonarr", 2) + if err != nil || total != 5 || len(list) != 2 || list[0].Sequence != 5 { + t.Fatalf("%d of %d (%v): %+v", len(list), total, err, list) + } + if _, total, _ := DeadLetters(js.Context(), "another_one", 2); total != 0 { + t.Fatalf("another consumer's total is %d", total) } } @@ -262,3 +309,28 @@ func TestTheControllerActsOnAnEventDeliveredAgain(t *testing.T) { } eventually(t, "the upgrade delivered again reaching the controller", func() bool { return told.count() == 1 }) } + +// The notices cannot be taken (their stream is gone): said as a condition, and the controller serves on. +func TestNoticesThatCannotBeTakenAreSaidAndServingGoesOn(t *testing.T) { + js := aBus(t) + if err := js.Context().DeleteStream(broker.DeadLetterNoticesStream); err != nil { + t.Fatal(err) + } + stop := keepingGivenUp(context.Background(), js, quiet()) + defer stop() + eventually(t, "the failure said", func() bool { + for _, a := range Advisories.Since(time.Now().Add(-time.Minute)) { + if a.Token == AdvisoryNotKept && a.Stream == broker.DeadLetterNoticesStream { + return true + } + } + return false + }) + held := &counted{} + _, stopServing := servingOn(t, js, held) + defer stopServing() + if _, err := js.Context().Publish("mesh.control.anchor.report", []byte(`{"node":"anchor"}`)); err != nil { + t.Fatal(err) + } + eventually(t, "a report heard while the notices cannot be taken", func() bool { return held.count() == 1 }) +} diff --git a/internal/link/receive_nats.go b/internal/link/receive_nats.go index 14f6eb16..e8ad09f3 100644 --- a/internal/link/receive_nats.go +++ b/internal/link/receive_nats.go @@ -41,6 +41,15 @@ type natsInbound struct { // that restarts loses these and starts the window again, which is correct — it is holding // nothing, and the messages are all still on the server. since map[uint64]time.Time + // log is where it says what it could not do; the standard logger when nobody gave one. + log *log.Logger +} + +func (n *natsInbound) logger() *log.Logger { + if n.log != nil { + return n.log + } + return log.Default() } // Nats is the consume side of the bus being built. @@ -70,7 +79,7 @@ func (n *natsInbound) Close() {} // which message is held, and since when — is read and written without a lock because the AMQP loop // never had two. A second goroutine would make that wrong in a way no test would catch. func (n *natsInbound) Receive(ctx context.Context, act func(context.Context, Control)) error { - if err := broker.AssertMeshConsumers(n.js); err != nil { + if err := broker.AssertServingConsumers(n.js); err != nil { return err } js, conn := n.js.Context(), n.js.Conn() @@ -96,11 +105,10 @@ func (n *natsInbound) Receive(ctx context.Context, act func(context.Context, Con // What a consumer gave up on, kept by the controller acting now (novox/hq issue 330): the server's // notices wait in their stream for it, so one said while no controller listened is not lost. - kept, err := keepGivenUp(js, log.Default()) - if err != nil { - return fmt.Errorf("taking the notices of messages consumers gave up on: %w", err) - } - defer func() { _ = kept.Unsubscribe() }() + // **Never the reason the controller stops serving**: a notice it cannot take yet waits in its stream, + // and the failure is said and tried again every minute. + stopKeeping := keepingGivenUp(ctx, n.js, n.logger()) + defer stopKeeping() // Heartbeats, on core NATS and off any stream (design 25 §3). Their own subscription because // they are their own guarantee: a lost one is the next one.