Files
mesh-controller/internal/link/bus.go
jschoubben 88bef39952 The consume side on NATS, and the window held by the server
The other implementation behind the seam, so the store-window guarantee now has
both: one loop, one message at a time, the same window deciding. What differs is
where a held message lives, and that is the whole point of the move — the AMQP
side keeps an unacknowledged delivery in this process, bounded by the prefetch
and lost if the controller stops; this keeps eight bytes saying when the window
opened, and the message stays the server's.

Checked against a running server, seven claims that reasoning cannot answer: a
report is heard and leaves the work queue; one the store cannot take is naked
with a delay, stays in the stream, and is recorded when the store returns; one
about a superseded declaration is settled without being acted on; one the store
never takes is let go once the bound passes; a heartbeat is heard and nothing is
persisted; and the enrolment answer reaches the address the request carried in
its payload — the test design 25 §2 asks for, so the reason for that field
cannot quietly become folklore.

Three things the wiring forced into the open:

**The controller could not have consumed a module event.** Its permissions
granted no event subject to subscribe and no ack subject on the events stream,
so every announcement would have been redelivered for ever, refused by the list
it already had. Both narrow: each followed subject named, not `mesh.mod.*.>`.

**The controller's consumers are not derived.** It files no manifest, so its
authority cannot come from a declaration that does not exist; they sit beside the
mesh's own streams and are asserted the same way. No max-deliver on CONTROL —
the window's bound is the controller's, and a server that dead-lettered first
would discard the push the stream exists to protect.

**Channels, not callbacks.** The library would run a handler on its own
goroutine, and the window's bookkeeping is unlocked because the AMQP loop never
had two.
2026-09-27 00:53:33 +02:00

195 lines
8.1 KiB
Go

package link
import (
"context"
"errors"
"fmt"
"time"
"github.com/nats-io/nats.go"
amqp "github.com/rabbitmq/amqp091-go"
)
// Bus is what the controller needs of the mesh's bus, **in the mesh's own words rather than a
// transport's** (novox/hq ADR 0116 step 3).
//
// Until now every one of these functions took the transport's own channel type, so the transport
// reached every caller and changing it meant touching all of them. The seam is small — the
// controller sends exactly two kinds of message that expect no answer, and asks two kinds of
// question — which is why the bus can be replaced at all.
//
// Two implementations live below, and both ship until the rollout (ADR 0116: nothing moves a
// node's bus before step 5). Both shipping is what makes them comparable — the same caller, the
// same arguments, and one conformance fixture holding them to one envelope.
type Bus interface {
// PublishEvent announces something that happened, under the emitter's own name. 1:many, and
// nobody is obliged to act (ADR 0041).
PublishEvent(ctx context.Context, key, source, node string, body []byte) error
// PublishDeclaration delivers one node what it should be. Addressed to that node alone: a
// declaration is not an event, and replaying yesterday's is actively harmful
// (design 29 §4, the *state* shape).
PublishDeclaration(ctx context.Context, node string, body []byte) error
// AskTool sends one question to a module's tool and awaits one answer. A tool nobody serves
// must say so **at once** rather than after the whole wait: the difference between "that
// module is down" and "that tool is slow" is the first thing a person asking wants.
AskTool(ctx context.Context, module, tool string, args []byte, timeout time.Duration) ([]byte, error)
}
// --- The bus the mesh runs on today -----------------------------------------------------
// OverCurrent is the bus the mesh runs on today, until the rollout.
type OverCurrent struct{ Channel *amqp.Channel }
func (b OverCurrent) PublishEvent(ctx context.Context, key, source, node string, body []byte) error {
id, err := eventID()
if err != nil {
return err
}
return b.Channel.PublishWithContext(ctx, EventsExchange, key, false, false, amqp.Publishing{
ContentType: "application/json",
DeliveryMode: amqp.Persistent,
MessageId: id,
Timestamp: time.Now().UTC(),
Body: body,
Headers: amqp.Table{
"x-event-id": id,
"x-source": source,
"x-node": node,
"x-time": time.Now().UTC().Format(time.RFC3339),
"content-type": "application/json",
},
})
}
// AskTool is implemented over the existing reply-queue machinery in ask.go; this seam does not
// change how it works today.
func (b OverCurrent) AskTool(ctx context.Context, module, tool string, args []byte, timeout time.Duration) ([]byte, error) {
answer, err := Ask(ctx, b.Channel, module, tool, args, timeout)
if err != nil {
return nil, err
}
return answer.Result, nil
}
func (b OverCurrent) PublishDeclaration(ctx context.Context, node string, body []byte) error {
// To the queue directly rather than through an exchange: a declaration is for one node, and
// routing it by name through a shared exchange would mean a binding per node that nothing
// removes when a node is retired.
return b.Channel.PublishWithContext(ctx, "", QueueFor(node), false, false, amqp.Publishing{
ContentType: "application/json",
DeliveryMode: amqp.Persistent,
Body: body,
})
}
// --- NATS, the bus being built ----------------------------------------------------------------
// OverNATS is the bus as a JetStream context.
type OverNATS struct {
Conn *nats.Conn
JS nats.JetStreamContext
}
// The subjects a node publishes on, and the controller listens to.
//
// One tree, and each name says who it is about: `mesh.control.<node>.…` is a node's own, which is
// what lets a node's account be granted exactly its own prefix and nothing of any other node's
// (design 25 §2, §4). The two that belong to no node — an enrolment, because a machine enrolling
// has no name the mesh has agreed to yet, and a build's outcome, because a builder is not
// reporting about itself — are named directly.
const (
// EnrolSubject is where a joining machine asks. Its enrolment user may publish here and
// nowhere else, so a leaked token buys nothing but the chance to enrol.
EnrolSubject = "mesh.control.enrol"
// BuiltSubject is where a build's outcome lands, for results nobody was waiting for.
BuiltSubject = "mesh.control.built"
// AliveSubjects is every node's heartbeat. Core NATS, never a stream: a lost heartbeat is the
// next heartbeat, and a stream of them is the mesh's least valuable message competing for
// retention with its most valuable (design 25 §3).
AliveSubjects = "mesh.control.*.alive"
)
// ReportSubject is where one node says what it did. On the CONTROL stream, because it is the
// message the store-window guarantee is about (ADR 0083).
func ReportSubject(node string) string { return "mesh.control." + node + ".report" }
// AliveSubject is one node's heartbeat.
func AliveSubject(node string) string { return "mesh.control." + node + ".alive" }
// EventSubject is where a module's event lands. Derived from the emitter, never taken from the
// caller: a source that could differ from the subject is an envelope that can lie about its
// origin, and on NATS the account's permissions make the subject the authority (design 29 §2).
func EventSubject(source, key string) string {
return "mesh.mod." + source + ".event." + key
}
// DeclareSubject is where one node's declaration lands. Last-per-subject on the NODES stream, so
// a node that was away gets exactly the current one and a replayed older one is refused by
// sequence — the wire-level answer to novox/hq issue 107.
func DeclareSubject(node string) string { return "mesh.node." + node + ".declare" }
func (b OverNATS) PublishEvent(ctx context.Context, key, source, node string, body []byte) error {
id, err := eventID()
if err != nil {
return err
}
h := nats.Header{}
h.Set("x-event-id", id)
h.Set("x-source", source)
h.Set("x-node", node)
h.Set("x-time", time.Now().UTC().Format(time.RFC3339))
h.Set("content-type", "application/json")
// The id is also the publish's message id, so the server refuses a duplicate inside its
// window. That narrows the window a consumer must deduplicate in; it does not remove the
// requirement, because the window is finite (design 19, delivery).
_, err = b.JS.PublishMsg(&nats.Msg{
Subject: EventSubject(source, key),
Header: h,
Data: body,
}, nats.MsgId(id), nats.Context(ctx))
if err != nil {
return fmt.Errorf("emitting %s: %w", key, err)
}
return nil
}
func (b OverNATS) PublishDeclaration(ctx context.Context, node string, body []byte) error {
_, err := b.JS.Publish(DeclareSubject(node), body, nats.Context(ctx))
if err != nil {
return fmt.Errorf("declaring to %s: %w", node, err)
}
return nil
}
// ToolSubject is where a module answers. Derived from the module and the tool, so a caller names
// what it wants rather than where it lives.
func ToolSubject(module, tool string) string { return "mesh.mod." + module + ".tool." + tool }
func (b OverNATS) AskTool(ctx context.Context, module, tool string, args []byte, timeout time.Duration) ([]byte, error) {
if len(args) == 0 {
args = []byte(`{}`)
}
ask, cancel := context.WithTimeout(ctx, timeout)
defer cancel()
// **No reply queue, and no correlation to check.** The caller's inbox is its own — each
// account is granted one prefix and no other (design 25 §4) — so an answer cannot reach the
// wrong asker and there is nothing to correlate against. That also settles a cost recorded
// in build.go: on a shared reply exchange every asker saw every result.
msg, err := b.Conn.RequestWithContext(ask, ToolSubject(module, tool), args)
if err != nil {
if errors.Is(err, nats.ErrNoResponders) {
// Said at once rather than after the whole wait: nothing is subscribed to that
// subject, which is a different fact from a tool being slow.
return nil, fmt.Errorf("nothing serves %s.%s", module, tool)
}
return nil, fmt.Errorf("asking %s.%s: %w", module, tool, err)
}
return msg.Data, nil
}