Deliver a dead letter again only where it is received, and never stop serving for the notices
mesh/delivery delivered
mesh/delivery-group group fix/330-a-message-given-up-on-is-kept delivered: every member is delivered
mesh/merge-gate pass: builds build-agent, mesh-controller, route-proxy → ace, g14, novox, shanks; no bus step; every machine composes with the change as it…
mesh/repo-check pass: its merge-check.sh passed

A dead letter was let go as delivered even when its consumer did not filter
its again subject; seat asks needed a grant over every seat's queue and left
the original stuck; a notices bind failure stopped the controller (review).
This commit is contained in:
jochen
2026-10-08 18:32:58 +02:00
parent 826dcb91b1
commit 909062e729
10 changed files with 236 additions and 53 deletions
+1 -1
View File
@@ -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
}
}
+4
View File
@@ -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")
}
+2
View File
@@ -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)
+3 -4
View File
@@ -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
+25
View File
@@ -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 {
+1 -1
View File
@@ -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" }
} }
+5 -2
View File
@@ -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,
}
}
+95 -25
View File
@@ -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 `<stream>.<consumer>`).
// 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 `<stream>.<consumer>`). 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.
+86 -14
View File
@@ -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 })
}
+14 -6
View File
@@ -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.