The controller's outbound link behind a seam, with both transports
Step 3.4, first half. Every one of these took an *amqp.Channel, so the transport reached every caller and swapping it meant touching all of them. The seam turned out to be small — the controller sends exactly two kinds of message that expect no answer — which is the same measurement that said this bus could be replaced at all. Bus is stated in the mesh's words, not a transport's: PublishEvent and PublishDeclaration. Two implementations, both shipping, because steps 1 to 4 leave every node on AMQP and the NATS one is selected at the rollout. Both ship is also what makes them comparable: one conformance fixture holds both to the same envelope, and the NATS one is checked against a real server reading back from the stream rather than from the code that wrote it. Still on *amqp.Channel: RequestBuild and Ask, which carry reply-queue machinery, and the whole consume side — the control loop, enrolment, serve.
This commit is contained in:
@@ -268,7 +268,7 @@ func answer(ctx context.Context, channel *amqp.Channel, publisher builder.Publis
|
|||||||
"manifest": json.RawMessage(result.Manifest), "against": result.Against,
|
"manifest": json.RawMessage(result.Manifest), "against": result.Against,
|
||||||
"made": result.Made,
|
"made": result.Made,
|
||||||
}
|
}
|
||||||
if err := link.EmitEvent(publishCtx, channel, link.KeyModuleBuilt, "builder", on, announced); err != nil {
|
if err := link.EmitEvent(publishCtx, link.OverAMQP{Channel: channel}, link.KeyModuleBuilt, "builder", on, announced); err != nil {
|
||||||
// Said, not fatal: the build happened and was answered. A module the catalogue has not
|
// Said, not fatal: the build happened and was answered. A module the catalogue has not
|
||||||
// heard of is a gap somebody can close; a build reported as failed because announcing
|
// heard of is a gap somebody can close; a build reported as failed because announcing
|
||||||
// it failed is a lie about work that was done.
|
// it failed is a lie about work that was done.
|
||||||
|
|||||||
@@ -142,7 +142,7 @@ func declare(ctx context.Context, args []string) error {
|
|||||||
}
|
}
|
||||||
defer server.Close()
|
defer server.Close()
|
||||||
|
|
||||||
if err := link.Declare(ctx, server.Channel(), ident, node, raw, 15*time.Second); err != nil {
|
if err := link.Declare(ctx, link.OverAMQP{Channel: server.Channel()}, ident, node, raw, 15*time.Second); err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
fmt.Printf("sent %s a signed declaration (%d bytes)\n", node, len(raw))
|
fmt.Printf("sent %s a signed declaration (%d bytes)\n", node, len(raw))
|
||||||
@@ -310,7 +310,7 @@ func pushCommand(ctx context.Context, args []string) error {
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
if err := link.Declare(ctx, server.Channel(), ident, s.node, body, 15*time.Second); err != nil {
|
if err := link.Declare(ctx, link.OverAMQP{Channel: server.Channel()}, ident, s.node, body, 15*time.Second); err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
// After it is away, not before. A digest recorded for something that failed to send would
|
// After it is away, not before. A digest recorded for something that failed to send would
|
||||||
@@ -393,7 +393,7 @@ func pushCommand(ctx context.Context, args []string) error {
|
|||||||
return declarationWith(held, open, node, plan, settings, gens, Allocating)
|
return declarationWith(held, open, node, plan, settings, gens, Allocating)
|
||||||
},
|
},
|
||||||
func(s readyNode, body []byte) error {
|
func(s readyNode, body []byte) error {
|
||||||
if err := link.Declare(ctx, server.Channel(), ident, s.node, body,
|
if err := link.Declare(ctx, link.OverAMQP{Channel: server.Channel()}, ident, s.node, body,
|
||||||
15*time.Second); err != nil {
|
15*time.Second); err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
@@ -611,7 +611,7 @@ func sendTo(ctx context.Context, open *stores, names []string) error {
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
if err := link.Declare(ctx, server.Channel(), ident, s.node, body, 15*time.Second); err != nil {
|
if err := link.Declare(ctx, link.OverAMQP{Channel: server.Channel()}, ident, s.node, body, 15*time.Second); err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
record, err := inv.NodeByName(ctx, s.node)
|
record, err := inv.NodeByName(ctx, s.node)
|
||||||
|
|||||||
@@ -0,0 +1,121 @@
|
|||||||
|
package link
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"fmt"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
amqp "github.com/rabbitmq/amqp091-go"
|
||||||
|
"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 an `*amqp.Channel`, so the transport reached every
|
||||||
|
// caller and swapping 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 could be replaced at all.
|
||||||
|
//
|
||||||
|
// Two implementations live below. Both ship: steps 1 to 4 leave every node on AMQP
|
||||||
|
// ([ADR 0116](novox/hq)), so the controller keeps speaking it and the NATS one is selected at the
|
||||||
|
// rollout. That is also what makes them comparable — the same caller, the same arguments, and a
|
||||||
|
// conformance fixture holding both 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
|
||||||
|
}
|
||||||
|
|
||||||
|
// --- AMQP, the bus the mesh runs on today -----------------------------------------------------
|
||||||
|
|
||||||
|
// OverAMQP is the bus as a channel.
|
||||||
|
type OverAMQP struct{ Channel *amqp.Channel }
|
||||||
|
|
||||||
|
func (b OverAMQP) 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",
|
||||||
|
},
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
func (b OverAMQP) 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{ JS nats.JetStreamContext }
|
||||||
|
|
||||||
|
// 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
|
||||||
|
}
|
||||||
@@ -0,0 +1,95 @@
|
|||||||
|
package link
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"encoding/json"
|
||||||
|
"os"
|
||||||
|
"testing"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"github.com/nats-io/nats.go"
|
||||||
|
)
|
||||||
|
|
||||||
|
// **Both implementations, one fixture.** The point of the seam is not that the transport can be
|
||||||
|
// swapped — it is that the two can be held to the same envelope while both ship, so the day the
|
||||||
|
// bus moves is a configuration change rather than a discovery.
|
||||||
|
//
|
||||||
|
// Against a real server, because what the fixture pins is what reaches the wire:
|
||||||
|
//
|
||||||
|
// docker run -d --rm --name t -p 14222:4222 nats:2.10-alpine -js
|
||||||
|
// MESH_TEST_NATS=nats://127.0.0.1:14222 go test ./internal/link/ -run TestTheNatsBus
|
||||||
|
func TestTheNatsBusEmitsTheEnvelopeTheFixturePins(t *testing.T) {
|
||||||
|
url := os.Getenv("MESH_TEST_NATS")
|
||||||
|
if url == "" {
|
||||||
|
t.Skip("MESH_TEST_NATS unset")
|
||||||
|
}
|
||||||
|
f := loadFixture(t, "events/module-event.json")
|
||||||
|
|
||||||
|
conn, err := nats.Connect(url)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
defer conn.Close()
|
||||||
|
js, err := conn.JetStream()
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if _, err := js.AddStream(&nats.StreamConfig{
|
||||||
|
Name: "EVENTS", Subjects: []string{"mesh.mod.*.event.>"},
|
||||||
|
}); err != nil && err != nats.ErrStreamNameAlreadyInUse {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
|
||||||
|
bus := OverNATS{JS: js}
|
||||||
|
body, _ := json.Marshal(f.Given.Body)
|
||||||
|
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
|
||||||
|
defer cancel()
|
||||||
|
if err := EmitEvent(ctx, bus, f.Given.Key, f.Given.Module, f.Given.Node, f.Given.Body); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
|
||||||
|
// Read back from the stream, not from the thing that wrote it.
|
||||||
|
raw, err := js.GetLastMsg("EVENTS", f.Wire.Subject)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("nothing landed on %s, which the fixture names: %v", f.Wire.Subject, err)
|
||||||
|
}
|
||||||
|
for _, h := range f.Wire.RequiredHeaders {
|
||||||
|
if raw.Header.Get(h) == "" {
|
||||||
|
t.Errorf("%s is not set on the wire, and the fixture requires it", h)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if got := raw.Header.Get("x-source"); got != f.Given.Module {
|
||||||
|
t.Errorf("x-source is %q; the subject says %q", got, f.Given.Module)
|
||||||
|
}
|
||||||
|
if string(raw.Data) != string(body) {
|
||||||
|
t.Errorf("the payload is %s, expected the body alone: %s", raw.Data, body)
|
||||||
|
}
|
||||||
|
// The envelope must not also be nested inside the payload.
|
||||||
|
var nested map[string]any
|
||||||
|
if json.Unmarshal(raw.Data, &nested) == nil {
|
||||||
|
if _, has := nested["key"]; has {
|
||||||
|
t.Error("the payload carries the envelope's own fields, which the fixture refuses")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// The subject a declaration lands on is one node's, and nothing else's — the state shape.
|
||||||
|
func TestADeclarationIsAddressedToOneNode(t *testing.T) {
|
||||||
|
if got := DeclareSubject("anchor"); got != "mesh.node.anchor.declare" {
|
||||||
|
t.Fatalf("a declaration would go to %q", got)
|
||||||
|
}
|
||||||
|
if DeclareSubject("anchor") == DeclareSubject("laptop") {
|
||||||
|
t.Fatal("two nodes share a declaration subject, so each would apply the other's")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// A module cannot emit under another's name: the subject is derived from the source, and the
|
||||||
|
// server's permissions make that subject the authority.
|
||||||
|
func TestAnEventsSubjectIsDerivedFromItsSource(t *testing.T) {
|
||||||
|
if got := EventSubject("shop", "order.placed"); got != "mesh.mod.shop.event.order.placed" {
|
||||||
|
t.Fatalf("an event would land on %q", got)
|
||||||
|
}
|
||||||
|
if EventSubject("shop", "x") == EventSubject("billing", "x") {
|
||||||
|
t.Fatal("two modules share an event subject, so neither owns its own name")
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -5,8 +5,6 @@ import (
|
|||||||
"encoding/json"
|
"encoding/json"
|
||||||
"fmt"
|
"fmt"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
amqp "github.com/rabbitmq/amqp091-go"
|
|
||||||
)
|
)
|
||||||
|
|
||||||
// Signer is whatever holds the control plane's signing key.
|
// Signer is whatever holds the control plane's signing key.
|
||||||
@@ -21,7 +19,7 @@ type Signer interface {
|
|||||||
// something else, and the node would refuse a declaration that was genuinely the mesh's.
|
// something else, and the node would refuse a declaration that was genuinely the mesh's.
|
||||||
//
|
//
|
||||||
// Published to the node's own queue, which its account alone may read.
|
// Published to the node's own queue, which its account alone may read.
|
||||||
func Declare(ctx context.Context, channel *amqp.Channel, signer Signer, node string,
|
func Declare(ctx context.Context, bus Bus, signer Signer, node string,
|
||||||
declaration []byte, timeout time.Duration) error {
|
declaration []byte, timeout time.Duration) error {
|
||||||
|
|
||||||
if !json.Valid(declaration) {
|
if !json.Valid(declaration) {
|
||||||
@@ -41,13 +39,5 @@ func Declare(ctx context.Context, channel *amqp.Channel, signer Signer, node str
|
|||||||
publish, cancel := context.WithTimeout(ctx, timeout)
|
publish, cancel := context.WithTimeout(ctx, timeout)
|
||||||
defer cancel()
|
defer cancel()
|
||||||
|
|
||||||
// Published to the queue directly rather than through the exchange: a declaration is for one
|
return bus.PublishDeclaration(publish, node, body)
|
||||||
// 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 channel.PublishWithContext(publish, "", QueueFor(node), false, false,
|
|
||||||
amqp.Publishing{
|
|
||||||
ContentType: "application/json",
|
|
||||||
DeliveryMode: amqp.Persistent,
|
|
||||||
Body: body,
|
|
||||||
})
|
|
||||||
}
|
}
|
||||||
|
|||||||
+7
-25
@@ -6,9 +6,6 @@ import (
|
|||||||
"encoding/hex"
|
"encoding/hex"
|
||||||
"encoding/json"
|
"encoding/json"
|
||||||
"fmt"
|
"fmt"
|
||||||
"time"
|
|
||||||
|
|
||||||
amqp "github.com/rabbitmq/amqp091-go"
|
|
||||||
)
|
)
|
||||||
|
|
||||||
// Emitting a module event from Go.
|
// Emitting a module event from Go.
|
||||||
@@ -29,32 +26,17 @@ const (
|
|||||||
|
|
||||||
// EmitEvent publishes one module event, in the envelope the sdk's consumers expect.
|
// EmitEvent publishes one module event, in the envelope the sdk's consumers expect.
|
||||||
//
|
//
|
||||||
// Persistent, because an event that a broker restart loses is not an announcement. The publish is
|
// The envelope is the transport's to write (bus.go) and this is only what goes in it, which is
|
||||||
// not confirmed here: the caller has already done the work the event describes, and a build that
|
// what lets one conformance fixture hold both implementations to the same headers.
|
||||||
// succeeded must not be reported as failed because saying so failed.
|
func EmitEvent(ctx context.Context, bus Bus, eventType, source, node string, body any) error {
|
||||||
func EmitEvent(ctx context.Context, channel *amqp.Channel, eventType, source, node string, body any) error {
|
|
||||||
payload, err := json.Marshal(body)
|
payload, err := json.Marshal(body)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return fmt.Errorf("cannot serialise a %s event: %w", eventType, err)
|
return fmt.Errorf("cannot serialise a %s event: %w", eventType, err)
|
||||||
}
|
}
|
||||||
id, err := eventID()
|
// The publish is not confirmed by the caller: it has already done the work the event
|
||||||
if err != nil {
|
// describes, and a build that succeeded must not be reported as failed because saying so
|
||||||
return err
|
// failed. Each transport decides what "published" means for it.
|
||||||
}
|
return bus.PublishEvent(ctx, eventType, source, node, payload)
|
||||||
return channel.PublishWithContext(ctx, EventsExchange, eventType, false, false, amqp.Publishing{
|
|
||||||
ContentType: "application/json",
|
|
||||||
DeliveryMode: amqp.Persistent,
|
|
||||||
MessageId: id,
|
|
||||||
Timestamp: time.Now().UTC(),
|
|
||||||
Body: payload,
|
|
||||||
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",
|
|
||||||
},
|
|
||||||
})
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// eventID is what a consumer deduplicates on: delivery is at-least-once, so a handler must be able
|
// eventID is what a consumer deduplicates on: delivery is at-least-once, so a handler must be able
|
||||||
|
|||||||
@@ -631,7 +631,7 @@ func (s *Server) catchingUp(ctx context.Context, delivery amqp.Delivery) {
|
|||||||
sent := 0
|
sent := 0
|
||||||
for _, a := range announcements {
|
for _, a := range announcements {
|
||||||
a.Replay = true
|
a.Replay = true
|
||||||
if err := EmitEvent(ctx, s.channel, KeyModuleBuilt, "control-plane", "", a); err != nil {
|
if err := EmitEvent(ctx, OverAMQP{Channel: s.channel}, KeyModuleBuilt, "control-plane", "", a); err != nil {
|
||||||
// Said and abandoned rather than retried: the catalogue asks again every time it
|
// Said and abandoned rather than retried: the catalogue asks again every time it
|
||||||
// starts, and half a graph delivered twice is no better than half delivered once.
|
// starts, and half a graph delivered twice is no better than half delivered once.
|
||||||
s.log.Printf("replaying %s at %s failed, and the rest is abandoned: %v",
|
s.log.Printf("replaying %s at %s failed, and the rest is abandoned: %v",
|
||||||
|
|||||||
Reference in New Issue
Block a user