package link import ( "context" "encoding/json" "errors" "fmt" "strings" "time" "github.com/nats-io/nats.go" "github.com/novox/mesh-controller/internal/broker" ) // The consume side on the bus being built. // // The shape is the AMQP one's, because the seam made them comparable: one loop, one message at a // time, and the same window deciding. What differs is where a held message lives — and that is the // whole point of the move. On the bus the mesh has, holding one means keeping an unacknowledged // delivery in this process, bounded by the prefetch and lost if the controller stops. Here it is a // `nak` with a delay: the message stays the server's, the controller keeps nothing but the moment // it first could not take it, and a controller that restarts mid-window has nothing to lose. // natsInbound consumes what nodes and modules say over NATS. type natsInbound struct { js *broker.JetStream // follows is the kinds asked for beyond what nodes say (Also). The events those are are the // only ones the controller subscribes, and only when something is listening. follows map[string]bool // since is when the controller first could not take a message, by that message's place in its // stream. // // **A timestamp, not a message.** This is the whole difference the move buys: the AMQP side // keeps the delivery, and this keeps eight bytes saying when the window opened. A controller // 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 } // Nats is the consume side of the bus being built. func Nats(js *broker.JetStream) Inbound { return &natsInbound{js: js, follows: map[string]bool{}, since: map[uint64]time.Time{}} } // Also records one more kind to subscribe. Nothing is subscribed here: the controller's consumer on // the events stream carries both of these as filters, so it is created once, in Receive, with // whatever was asked for — and not at all when nothing was. func (n *natsInbound) Also(kind string) error { switch kind { case KindModuleMoved, KindCatchUp: n.follows[kind] = true return nil default: return fmt.Errorf("nothing subscribes %s separately on this bus", kind) } } func (n *natsInbound) Close() {} // Receive consumes until the context ends. // // Three subscriptions, and each is a channel the one loop selects on. **Channels rather than // callbacks**: the library would run a handler on its own goroutine, and the window's bookkeeping — // 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 { return err } js, conn := n.js.Context(), n.js.Conn() // What nodes say, off the CONTROL stream. Bound to the durable the controller asserted rather // than creating one here: the consumer is an object with a configuration — ack policy, ack // wait, redelivery — and a client that creates its own would be a second opinion about it. control := make(chan *nats.Msg, Prefetch) said, err := js.ChanSubscribe("", control, nats.Bind("CONTROL", broker.ControllerName)) if err != nil { return fmt.Errorf("subscribing to what nodes say: %w", err) } defer func() { _ = said.Unsubscribe() }() // 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. beats := make(chan *nats.Msg, Prefetch) alive, err := conn.ChanSubscribe(AliveSubjects, beats) if err != nil { return fmt.Errorf("subscribing to heartbeats: %w", err) } defer func() { _ = alive.Unsubscribe() }() // The events the controller follows, when something is listening for them. var events chan *nats.Msg if len(n.follows) > 0 { events = make(chan *nats.Msg, Prefetch) followed, err := js.ChanSubscribe("", events, nats.Bind("EVENTS", broker.ControllerName)) if err != nil { return fmt.Errorf("subscribing to what the catalogue says: %w", err) } defer func() { _ = followed.Unsubscribe() }() } // A connection that dropped is said, not discovered. A controller whose bus connection is gone // is a mesh where nothing can be told anything. gone := make(chan error, 1) conn.SetDisconnectErrHandler(func(_ *nats.Conn, err error) { select { case gone <- err: default: } }) for { select { case <-ctx.Done(): return nil case err := <-gone: return fmt.Errorf("the bus connection dropped: %w", err) case msg := <-beats: n.deliver(ctx, act, msg, false) case msg := <-events: n.deliver(ctx, act, msg, true) case msg, ok := <-control: if !ok { return errors.New("the bus stopped delivering") } n.deliver(ctx, act, msg, true) } } } // deliver names one message and hands it to the loop, or drops it where the mesh has no name for // its subject — which cannot happen through a filter the controller wrote, and is said rather than // ignored for exactly that reason. func (n *natsInbound) deliver(ctx context.Context, act func(context.Context, Control), msg *nats.Msg, streamed bool) { kind, known := kindOfSubject(msg.Subject) if !known { if streamed { _ = msg.Term() } return } m := &natsControl{kind: kind, msg: msg, on: n} if streamed { // A message with no metadata is not from a stream, whatever it was delivered on, and the // window has nothing to hold it by. Said by leaving the sequence at zero. if meta, err := msg.Metadata(); err == nil { m.seq = meta.Sequence.Stream m.delivered = meta.NumDelivered } } act(ctx, m) } // kindOfSubject is how this transport's addressing becomes what the mesh calls a message. // // By subject, which is the only thing the server enforces: a body claiming to be a report does not // make it one, and on this bus the subject an account may publish *is* its authority (design 29 // §2). The mirror of the routing-key table on the bus the mesh has. func kindOfSubject(subject string) (string, bool) { switch subject { case EnrolSubject: return KindEnrolment, true case BuiltSubject: return KindBuilt, true } if node, rest, ok := strings.Cut(strings.TrimPrefix(subject, "mesh.control."), "."); ok && node != "" && !strings.Contains(node, ".") { switch rest { case "report": return KindReport, true case "alive": return KindHeartbeat, true } } switch subject { case broker.ControllerFollows[0]: return KindModuleMoved, true case broker.ControllerFollows[1]: return KindCatchUp, true } return "", false } // natsControl is one message from the bus being built, as the controller reads it. type natsControl struct { kind string msg *nats.Msg on *natsInbound // seq is this message's place in its stream; zero for a core message, which has none and // cannot be held. seq uint64 // delivered is how many times the server has handed this message over, this time included. delivered uint64 } func (m *natsControl) Kind() string { return m.kind } func (m *natsControl) Body() []byte { return m.msg.Data } // Redelivered is what the server counted, not what the controller remembers. Which is the answer to // a question the AMQP side could only guess at across a restart: an enrolment redelivered because // the controller stopped mid-answer reads as redelivered to the controller that comes back. func (m *natsControl) Redelivered() bool { return m.delivered > 1 } func (m *natsControl) HeldFor() time.Duration { if m.seq == 0 { return 0 } first, held := m.on.since[m.seq] if !held { return 0 } return time.Since(first) } // About is nothing here, and that is the point. // // Setting a held message aside when a newer one about the same thing arrives is what a controller // holding deliveries in memory can do. A naked message belongs to the server and comes back // whatever happened meanwhile, so the question "is this the past?" is answered by what the message // says instead — the digest of the declaration a report is about (window.go, design 25 §3). func (m *natsControl) About(string) {} // Answer publishes to the reply subject the request carries **in its payload**. // // Not `Respond`, and not the message's reply field: a message a JetStream consumer delivers has had // that field claimed for the consumer's own ack address, so answering it would send the reply to // `$JS.ACK.CONTROL.controller.…` and the enrolling node would wait out its timeout. Verified // against a running server (design 25 §2), which is why it is a field of the request and this reads // it from there. func (m *natsControl) Answer(ctx context.Context, body []byte) error { var addressed replyAddressed if err := json.Unmarshal(m.msg.Data, &addressed); err != nil { return fmt.Errorf("that request cannot be read, so its reply address cannot be: %w", err) } if addressed.ReplyTo == "" { return errors.New("that request named no reply subject in its payload, so nothing can be " + "told the answer") } return m.on.js.Conn().PublishMsg(&nats.Msg{Subject: addressed.ReplyTo, Data: body}) } func (m *natsControl) Took() error { m.forget() if m.seq == 0 { // Core NATS: nothing is keeping it, so there is nothing to settle. return nil } return m.msg.Ack(nats.Context(context.Background())) } // Drop terminates the delivery: understood, and the server is told not to send it again. Different // from an ack only in the server's own accounting, which is where somebody asking "what happened to // that message" will look. func (m *natsControl) Drop() error { m.forget() if m.seq == 0 { return nil } return m.msg.Term() } // Hold hands the message back with a delay, and remembers when the window opened. func (m *natsControl) Hold(after time.Duration) error { if m.seq == 0 { return errors.New("a message that is not in a stream cannot be held: nothing is keeping it") } if _, already := m.on.since[m.seq]; !already { m.on.since[m.seq] = time.Now() } return m.msg.NakWithDelay(after) } func (m *natsControl) forget() { if m.on != nil && m.seq != 0 { delete(m.on.since, m.seq) } } // replyAddressed is the one field every message that expects an answer carries. type replyAddressed struct { ReplyTo string `json:"reply_to,omitempty"` }