package link import ( "context" "errors" "fmt" "time" "github.com/nats-io/nats.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 // PublishSeatEvent states a fact under a role's own name, for the holder of that role. A // module's event is addressed to the module; a role's is addressed to the role, so it keeps // meaning when the holder changes (novox/hq ADR 0121, ADR 0129). PublishSeatEvent(ctx context.Context, seat, event 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 ----------------------------------------------------- // --- 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..…` 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 } // SeatEventSubject is where a role's own event lands. Derived from the role, never from its holder: // a fact about the build machine or about the control plane keeps its address when the module holding // that role is replaced (novox/hq ADR 0121, ADR 0129). func SeatEventSubject(seat, event string) string { return "mesh.seat." + seat + ".event." + event } // 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 } // PublishSeatEvent states a role's own fact. Same envelope as a module's event and a different // address: the source header is the role, because that is what the fact is about. func (b OverNATS) PublishSeatEvent(ctx context.Context, seat, event 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", seat) _, err = b.JS.PublishMsg(&nats.Msg{ Subject: SeatEventSubject(seat, event), Header: h, Data: body, }, nats.MsgId(id), nats.Context(ctx)) if err != nil { return fmt.Errorf("stating %s of the %s seat: %w", event, seat, 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 }