diff --git a/Makefile b/Makefile index 19a039f..6c69db8 100644 --- a/Makefile +++ b/Makefile @@ -21,10 +21,14 @@ export MESH_TEST_POSTGRES ?= postgres://postgres:check@127.0.0.1:$(PG_PORT)/post build: CGO_ENABLED=0 go build -trimpath -ldflags '$(LDFLAGS)' -o build/mesh-control ./cmd/mesh-control +# Tagged 'development' as well as by version, because the lab places images by name and a +# scenario naming a version would have to be edited on every build. The version tag is what a +# real bundle pins. IMAGE ?= mesh-control:$(VERSION) +DEV_TAG ?= mesh-control:development image: - docker build --build-arg VERSION=$(VERSION) -t $(IMAGE) . + docker build --build-arg VERSION=$(VERSION) -t $(IMAGE) -t $(DEV_TAG) . @echo @docker image inspect $(IMAGE) --format 'built {{.RepoTags}} {{.Size}} bytes' diff --git a/cmd/mesh-control/main.go b/cmd/mesh-control/main.go index 63f1d8b..c5938a6 100644 --- a/cmd/mesh-control/main.go +++ b/cmd/mesh-control/main.go @@ -19,6 +19,7 @@ import ( "github.com/novox/mesh-control/internal/broker" "github.com/novox/mesh-control/internal/identity" "github.com/novox/mesh-control/internal/inventory" + "github.com/novox/mesh-control/internal/link" "github.com/novox/mesh-control/internal/store" "github.com/novox/mesh-control/internal/token" ) @@ -67,6 +68,8 @@ func run() error { return identityCommand(ctx, args[1:]) case "broker": return brokerCommand(args[1:]) + case "serve": + return serve(ctx) case "version": fmt.Println(version) return nil @@ -89,6 +92,7 @@ func usage() { token issue --new create the record and issue for it identity show this control plane's signing key broker show where the broker is, and what to expect there + serve consume what nodes say, and answer version what this binary is Each context reaches its own store through its own credential (novox/hq ADR 0008), named @@ -253,6 +257,20 @@ func tokenCommand(ctx context.Context, args []string) error { return err } + // The account is created before the token is handed over, which is what removes the + // chicken-and-egg entirely: the mesh runs the broker, so a joining node's credentials can + // exist before it does. The one-time secret IS the password, so a node's first connection is + // already authenticated and enrolment is what happens over it. + if management, err := broker.ManagementFromEnvironment(); err == nil { + if err := management.CreateNodeAccount(ctx, issued.Node.Name, issued.Secret); err != nil { + return err + } + fmt.Printf("broker account %s created, scoped to %s and the %s exchange\n\n", + issued.Node.Name, link.QueueFor(issued.Node.Name), link.Exchange) + } else if !errors.Is(err, broker.ErrNotConfigured) { + return err + } + made := token.Token{Signer: key.Public, Secret: issued.Secret} // Absent is a state, not a failure: a control plane can hold records and a key before it has @@ -342,3 +360,40 @@ func brokerCommand(args []string) error { "node checks it before sending anything (novox/hq ADR 0004).\n") return nil } + +// serve is the control plane running: one connection to the broker, one queue, one consumer. +func serve(ctx context.Context) error { + inv, err := openInventory(ctx) + if err != nil { + return err + } + defer inv.Close() + + ident, err := openIdentity(ctx) + if err != nil { + return err + } + defer ident.Close() + + // Established at start rather than on first use. A control plane that cannot sign is one + // whose declarations every node correctly refuses, and that should be a startup failure + // rather than something discovered at the first declaration. + key, err := ident.Establish(ctx) + if err != nil { + return err + } + fmt.Printf("signing as %s\n", key.Fingerprint()[:16]) + + management, err := broker.ManagementFromEnvironment() + if err != nil && !errors.Is(err, broker.ErrNotConfigured) { + return err + } + + server, err := link.Connect(link.Enrolment{Inventory: inv, Identity: ident, Broker: management}) + if err != nil { + return err + } + defer server.Close() + + return server.Serve(ctx) +} diff --git a/go.mod b/go.mod index 5f66f07..c8188a4 100644 --- a/go.mod +++ b/go.mod @@ -7,6 +7,7 @@ require ( github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761 // indirect github.com/jackc/pgx/v5 v5.10.0 // indirect github.com/jackc/puddle/v2 v2.2.2 // indirect + github.com/rabbitmq/amqp091-go v1.14.0 // indirect golang.org/x/sync v0.17.0 // indirect golang.org/x/text v0.29.0 // indirect ) diff --git a/go.sum b/go.sum index fc43d31..6806b42 100644 --- a/go.sum +++ b/go.sum @@ -8,6 +8,8 @@ github.com/jackc/pgx/v5 v5.10.0/go.mod h1:mal1tBGAFfLHvZzaYh77YS/eC6IX9OWbRV1QII github.com/jackc/puddle/v2 v2.2.2 h1:PR8nw+E/1w0GLuRFSmiioY6UooMp6KJv0/61nB7icHo= github.com/jackc/puddle/v2 v2.2.2/go.mod h1:vriiEXHvEE654aYKXXjOvZM39qJ0q+azkZFrfEOc3H4= github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= +github.com/rabbitmq/amqp091-go v1.14.0 h1:RSaT7aOKt/OrkVUyswPDW29lnRz9psuGmfZFBmLqLek= +github.com/rabbitmq/amqp091-go v1.14.0/go.mod h1:Hy4jKW5kQART1u+JkDTF9YYOQUHXqMuhrgxOEeS7G4o= github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME= github.com/stretchr/testify v1.3.0/go.mod h1:M5WIy9Dh21IEIfnGCwXGc5bZfKNJtfHm1UVUgZn+9EI= github.com/stretchr/testify v1.7.0/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= diff --git a/internal/broker/management.go b/internal/broker/management.go new file mode 100644 index 0000000..c2b3cc0 --- /dev/null +++ b/internal/broker/management.go @@ -0,0 +1,187 @@ +package broker + +import ( + "bytes" + "context" + "encoding/json" + "fmt" + "io" + "net/http" + "net/url" + "os" + "regexp" + "strings" + "time" +) + +// The mesh runs the broker, so there is no chicken-and-egg in a node needing an account before it +// can connect: the account is created when the token is issued, and the one-time secret in that +// token IS the password. A node's first connection is already authenticated, and enrolment is +// what happens over it. +// +// novox/hq ADR 0004's *a node holds its own identity and nothing else* is why the account is per +// node rather than shared. A shared enrolment account would let any node consume another's queue, +// which is the shared-credential fault that record exists to remove, reappearing at the transport. + +// ManagementVar holds the broker's management API, credentials included. +const ManagementVar = "MESH_BROKER_MANAGEMENT" + +// safeName is what a node may be called at the broker. +// +// The name goes into a URL path and into permission patterns, which are regular expressions. A +// name carrying a `.` or a `*` would silently widen what that node may reach — so it is +// constrained here rather than escaped later, because an escape that is forgotten once is a node +// reading everybody's queues. +var safeName = regexp.MustCompile(`^[a-z0-9][a-z0-9-]{0,62}$`) + +// Management is the broker's administrative interface. +type Management struct { + base *url.URL + client *http.Client +} + +// ManagementFromEnvironment reads where the management API is, if it is configured. +func ManagementFromEnvironment() (*Management, error) { + raw := strings.TrimSpace(os.Getenv(ManagementVar)) + if raw == "" { + return nil, ErrNotConfigured + } + base, err := url.Parse(raw) + if err != nil || base.Host == "" { + // The value carries a password, so it is not quoted back. + return nil, fmt.Errorf("%s is not a usable URL", ManagementVar) + } + return &Management{base: base, client: &http.Client{Timeout: 15 * time.Second}}, nil +} + +// QueueFor is the queue a node consumes from. One per node, named after it. +func QueueFor(node string) string { return "node." + node } + +// ExchangeName is where nodes publish what they have to say. One exchange, and the control plane +// is the only consumer behind it (novox/hq ADR 0006 — one consumer, so two cannot silently split +// the traffic between them). +const ExchangeName = "mesh" + +// CreateNodeAccount gives a node its own broker account, with the token's secret as the password. +// +// Scoped so a node can reach its own queue and the one exchange, and nothing else. The patterns +// are anchored: a node called `laptop` must not be able to read `laptop-of-somebody-else`. +func (m *Management) CreateNodeAccount(ctx context.Context, node, password string) error { + if !safeName.MatchString(node) { + return fmt.Errorf( + "%q cannot be a broker account name: it becomes part of a permission pattern, so it "+ + "is lower-case letters, digits and dashes", node) + } + + if err := m.put(ctx, "/api/users/"+url.PathEscape(node), + map[string]string{"password": password, "tags": ""}); err != nil { + return fmt.Errorf("cannot create the broker account for %s: %w", node, err) + } + + queue := regexp.QuoteMeta(QueueFor(node)) + if err := m.put(ctx, "/api/permissions/%2f/"+url.PathEscape(node), map[string]string{ + "configure": "^" + queue + "$", + "write": "^(" + regexp.QuoteMeta(ExchangeName) + "|" + queue + ")$", + "read": "^" + queue + "$", + }); err != nil { + return fmt.Errorf("cannot scope the broker account for %s: %w", node, err) + } + return nil +} + +// RemoveNodeAccount withdraws a node's access. +func (m *Management) RemoveNodeAccount(ctx context.Context, node string) error { + if !safeName.MatchString(node) { + return fmt.Errorf("%q is not a broker account name", node) + } + return m.do(ctx, http.MethodDelete, "/api/users/"+url.PathEscape(node), nil) +} + +// Accounts lists the broker's users, so a picture can be read from the system rather than assumed +// (novox/hq ADR 0018). +func (m *Management) Accounts(ctx context.Context) ([]string, error) { + body, err := m.get(ctx, "/api/users") + if err != nil { + return nil, err + } + var users []struct { + Name string `json:"name"` + } + if err := json.Unmarshal(body, &users); err != nil { + return nil, err + } + names := make([]string, 0, len(users)) + for _, u := range users { + names = append(names, u.Name) + } + return names, nil +} + +func (m *Management) put(ctx context.Context, path string, body any) error { + return m.do(ctx, http.MethodPut, path, body) +} + +func (m *Management) get(ctx context.Context, path string) ([]byte, error) { + request, err := m.request(ctx, http.MethodGet, path, nil) + if err != nil { + return nil, err + } + response, err := m.client.Do(request) + if err != nil { + return nil, err + } + defer response.Body.Close() + if response.StatusCode >= 300 { + return nil, fmt.Errorf("the broker's management API answered %s to GET %s", + response.Status, path) + } + return io.ReadAll(io.LimitReader(response.Body, 1<<20)) +} + +func (m *Management) do(ctx context.Context, method, path string, body any) error { + request, err := m.request(ctx, method, path, body) + if err != nil { + return err + } + response, err := m.client.Do(request) + if err != nil { + return err + } + defer response.Body.Close() + if response.StatusCode >= 300 { + detail, _ := io.ReadAll(io.LimitReader(response.Body, 4096)) + return fmt.Errorf("the broker's management API answered %s to %s %s: %s", + response.Status, method, path, strings.TrimSpace(string(detail))) + } + return nil +} + +func (m *Management) request(ctx context.Context, method, path string, body any) (*http.Request, error) { + var payload io.Reader + if body != nil { + raw, err := json.Marshal(body) + if err != nil { + return nil, err + } + payload = bytes.NewReader(raw) + } + + // Path joined by hand rather than through url.Parse: %2f is the default vhost and must reach + // the broker still encoded. Parsing would decode it to a slash and address a different route. + target := strings.TrimSuffix(m.base.String(), "/") + if user := m.base.User; user != nil { + target = strings.TrimSuffix(m.base.Scheme+"://"+m.base.Host, "/") + } + request, err := http.NewRequestWithContext(ctx, method, target+path, payload) + if err != nil { + return nil, err + } + if user := m.base.User; user != nil { + password, _ := user.Password() + request.SetBasicAuth(user.Username(), password) + } + if body != nil { + request.Header.Set("Content-Type", "application/json") + } + return request, nil +} diff --git a/internal/inventory/nodes.go b/internal/inventory/nodes.go index b3c7e0f..ee69991 100644 --- a/internal/inventory/nodes.go +++ b/internal/inventory/nodes.go @@ -6,6 +6,7 @@ import ( "crypto/sha256" "encoding/base64" "encoding/hex" + "encoding/json" "errors" "fmt" "strings" @@ -208,3 +209,27 @@ func (i *Inventory) Redeem(ctx context.Context, secret string) (Node, error) { `select id, name, created from node where id = $1`, id).Scan(&n.ID, &n.Name, &n.Created) return n, err } + +// RecordProfile keeps the last thing a node said about what it can do. +// +// The last one, not a history: the control plane needs to know what this machine can run *now* in +// order to decide what it should run, and an old profile is worse than none — it describes a +// machine that may have been rebuilt since. +func (i *Inventory) RecordProfile(ctx context.Context, node string, profile map[string]any) error { + raw, err := json.Marshal(profile) + if err != nil { + return err + } + _, err = i.store.Pool().Exec(ctx, + `update node set profile = $2, last_seen = now() where id = $1`, node, raw) + return err +} + +// Seen records that a node was heard from. +// +// Separate from the profile because it happens far more often: a node reports it is alive +// constantly and describes itself rarely. +func (i *Inventory) Seen(ctx context.Context, node string) error { + _, err := i.store.Pool().Exec(ctx, `update node set last_seen = now() where id = $1`, node) + return err +} diff --git a/internal/link/enrolment.go b/internal/link/enrolment.go new file mode 100644 index 0000000..b58f93e --- /dev/null +++ b/internal/link/enrolment.go @@ -0,0 +1,67 @@ +package link + +import ( + "context" + "crypto/ed25519" + "errors" + "fmt" + + "github.com/novox/mesh-control/internal/broker" + "github.com/novox/mesh-control/internal/identity" + "github.com/novox/mesh-control/internal/inventory" +) + +// Enrolment is what actually happens when a node presents a token: the token is spent, the key is +// recorded, and the node gets its own queue. +// +// It reaches across two contexts and reads neither one's store from the other (novox/hq ADR 0008) +// — it holds both grants and asks each for its part, which is what the process running them is +// for. +type Enrolment struct { + Inventory *inventory.Inventory + Identity *identity.Identity + Broker *broker.Management +} + +// Enrol spends the token and records what the node presented. +// +// Order matters and it is the order things become irreversible. The token is spent first, in a +// single statement that both finds and marks it, so two machines racing on one secret produce one +// winner. Only then is a key recorded — because recording a key for a node whose token turned out +// to be spent would leave the mesh believing a machine that never had the right to join. +func (e Enrolment) Enrol(ctx context.Context, secret string, public ed25519.PublicKey, + profile map[string]any) (string, error) { + + if len(public) != ed25519.PublicKeySize { + return "", fmt.Errorf("a node presented a %d-byte key, and an identity is %d", + len(public), ed25519.PublicKeySize) + } + + node, err := e.Inventory.Redeem(ctx, secret) + if err != nil { + return "", err + } + + // From here the token is gone whatever happens next, so anything that fails leaves a node + // record with no live key — which is visible and fixable with a new token, where a spent + // token believed to be unspent is neither. + if _, err := e.Identity.RecordNodeKey(ctx, node.ID, public); err != nil { + return "", fmt.Errorf("the token was spent and the key could not be recorded, so %s has "+ + "no identity and needs a new token: %w", node.Name, err) + } + + if profile != nil { + if err := e.Inventory.RecordProfile(ctx, node.ID, profile); err != nil { + // Not fatal. The profile is what the control plane needs in order to decide what this + // machine should run, and it is reported again on every connection — so losing it + // here costs a decision that can be made later, not the enrolment. + return node.Name, nil + } + } + return node.Name, nil +} + +var _ Enroller = Enrolment{} + +// ErrNoBrokerManagement is returned when an account cannot be made because nothing was configured. +var ErrNoBrokerManagement = errors.New("no broker management configured") diff --git a/internal/link/protocol.go b/internal/link/protocol.go new file mode 100644 index 0000000..337253a --- /dev/null +++ b/internal/link/protocol.go @@ -0,0 +1,63 @@ +// Package link is the control plane's side of the connection nodes hold open. +// +// novox/hq ADR 0002: nodes communicate over a message broker, not over HTTP. One exchange, and +// the control plane is the single consumer behind it — ADR 0006 makes that a property worth +// having rather than an accident, because two consumers sharing a queue silently split the +// traffic between them, each receiving half of what it expects. That has happened here before. +package link + +// Exchange is where nodes publish everything they have to say. +const Exchange = "mesh" + +// ControlQueue is what the control plane consumes. One queue, one consumer. +const ControlQueue = "control" + +// Routing keys. A node may publish these; it may not publish anything else, because its broker +// account is scoped to this exchange and its own queue. +const ( + KeyEnrol = "enrol" +) + +// QueueFor is the queue a node consumes from — the only one it may read. +func QueueFor(node string) string { return "node." + node } + +// EnrolRequest is what a joining node says. +// +// It arrives on a connection the broker has already authenticated, because the account was +// created when the token was issued and the token's secret is its password. So this message is +// not how a node gets in — it is what it says once it is in. +type EnrolRequest struct { + // Node is what this machine believes it is called. Checked against the token, never trusted. + Node string `json:"node"` + + // Secret is the one-time right to join. The account password and this are the same string, + // which is deliberate: the broker proves somebody holds the token, and this proves the same + // thing to the control plane without the control plane having to ask the broker who connected. + Secret string `json:"secret"` + + // PublicKey is what the mesh will believe from now on. The node generated it; the private + // half has never left that machine (novox/hq ADR 0004). + PublicKey []byte `json:"public_key"` + + // Profile is what this machine can be asked to do. The control plane cannot decide what a + // node should run without it, so it arrives with enrolment rather than being asked for after. + Profile map[string]any `json:"profile,omitempty"` +} + +// EnrolReply is what the mesh says back. +type EnrolReply struct { + // Accepted says whether the node is now known. + Accepted bool `json:"accepted"` + + // Node is the name the mesh has for this machine, which settles any disagreement: the token + // was issued for a node record, and that record's name wins over what the machine called + // itself. + Node string `json:"node,omitempty"` + + // Queue is where this node listens from now on. + Queue string `json:"queue,omitempty"` + + // Refusal says why not, in words for a person. Deliberately the same for every reason a + // token can fail — unknown, spent, expired — so that guessing learns nothing. + Refusal string `json:"refusal,omitempty"` +} diff --git a/internal/link/serve.go b/internal/link/serve.go new file mode 100644 index 0000000..43c6afa --- /dev/null +++ b/internal/link/serve.go @@ -0,0 +1,189 @@ +package link + +import ( + "context" + "crypto/ed25519" + "encoding/json" + "errors" + "fmt" + "log" + "os" + "strings" + "time" + + amqp "github.com/rabbitmq/amqp091-go" +) + +// AMQPVar is the control plane's own connection to the broker. +const AMQPVar = "MESH_BROKER_AMQP" + +// Enroller is what the control plane does with an enrolment request. +// +// An interface so the serving loop can be tested against a real broker without a database, and +// so the two concerns — moving messages, and deciding — stay apart. +type Enroller interface { + // Enrol spends the token, records the key, and reports the node's name. The error is + // returned to the node as a refusal; it must be the same for every reason a token can fail. + Enrol(ctx context.Context, secret string, public ed25519.PublicKey, profile map[string]any) (string, error) +} + +// Server consumes what nodes say. +type Server struct { + conn *amqp.Connection + channel *amqp.Channel + enroller Enroller + log *log.Logger +} + +// Connect opens the control plane's own connection to the broker. +func Connect(enroller Enroller) (*Server, error) { + url := strings.TrimSpace(os.Getenv(AMQPVar)) + if url == "" { + return nil, fmt.Errorf( + "this control plane has no %s, so it cannot reach its broker. Nodes talk to it over "+ + "the broker and nowhere else, so without this it can hold records and answer "+ + "nothing", AMQPVar) + } + + conn, err := amqp.Dial(url) + if err != nil { + // Not quoted back: the URL carries the control plane's own broker password. + return nil, fmt.Errorf("cannot reach the broker named in %s: %w", AMQPVar, err) + } + channel, err := conn.Channel() + if err != nil { + conn.Close() + return nil, err + } + + // Declared here rather than assumed. The control plane is the only thing that may create + // them — a node's account can write to this exchange and read its own queue, and configure + // nothing else, so a node arriving before the control plane has ever run finds nothing and + // says so, rather than quietly creating a topology nobody designed. + if err := channel.ExchangeDeclare(Exchange, "direct", true, false, false, false, nil); err != nil { + conn.Close() + return nil, fmt.Errorf("cannot declare the %s exchange: %w", Exchange, err) + } + if _, err := channel.QueueDeclare(ControlQueue, true, false, false, false, nil); err != nil { + conn.Close() + return nil, fmt.Errorf("cannot declare the %s queue: %w", ControlQueue, err) + } + if err := channel.QueueBind(ControlQueue, KeyEnrol, Exchange, false, nil); err != nil { + conn.Close() + return nil, err + } + + return &Server{conn: conn, channel: channel, enroller: enroller, + log: log.New(os.Stdout, "", log.LstdFlags)}, nil +} + +func (s *Server) Close() { + if s.channel != nil { + _ = s.channel.Close() + } + if s.conn != nil { + _ = s.conn.Close() + } +} + +// Serve consumes until the context ends. +// +// One consumer, deliberately: with two, the broker would round-robin between them and each would +// receive half of what it expects — a fault this project has already had, between a module's +// daemon and its capability server. +func (s *Server) Serve(ctx context.Context) error { + // Prefetch of one. The control plane writes to a database per message, and a burst of + // enrolments delivered all at once would be held in memory rather than left on the broker, + // which is the one place they survive a restart. + if err := s.channel.Qos(1, 0, false); err != nil { + return err + } + + deliveries, err := s.channel.ConsumeWithContext(ctx, ControlQueue, "control-plane", + false, false, false, false, nil) + if err != nil { + return err + } + + closed := s.conn.NotifyClose(make(chan *amqp.Error, 1)) + s.log.Printf("consuming %s, bound to %s/%s", ControlQueue, Exchange, KeyEnrol) + + for { + select { + case <-ctx.Done(): + return nil + case reason := <-closed: + // Said rather than returned quietly. A control plane whose broker connection dropped + // is a mesh where nothing can be told anything, and the reason is the first thing + // anybody will want. + return fmt.Errorf("the broker connection closed: %v", reason) + case delivery, ok := <-deliveries: + if !ok { + return errors.New("the broker stopped delivering") + } + s.handle(ctx, delivery) + } + } +} + +func (s *Server) handle(ctx context.Context, delivery amqp.Delivery) { + switch delivery.RoutingKey { + case KeyEnrol: + s.handleEnrol(ctx, delivery) + default: + // Rejected without requeue: a message nothing understands will not be understood on the + // next attempt either, and requeuing it would spin. + s.log.Printf("refusing a message with routing key %q", delivery.RoutingKey) + _ = delivery.Reject(false) + } +} + +func (s *Server) handleEnrol(ctx context.Context, delivery amqp.Delivery) { + reply := EnrolReply{Refusal: "that token cannot be used"} + + var request EnrolRequest + if err := json.Unmarshal(delivery.Body, &request); err != nil { + s.log.Printf("an enrolment request could not be read: %v", err) + } else { + name, err := s.enroller.Enrol(ctx, request.Secret, request.PublicKey, request.Profile) + if err != nil { + // Logged in full here, where an operator can see it; sent back as one refusal, so + // that somebody guessing learns nothing from which reason came back. + s.log.Printf("refusing enrolment for %q: %v", request.Node, err) + } else { + reply = EnrolReply{Accepted: true, Node: name, Queue: QueueFor(name)} + s.log.Printf("enrolled %s", name) + } + } + + s.reply(ctx, delivery, reply) + + // Acknowledged after the reply is sent, so a control plane that dies mid-answer leaves the + // request on the broker rather than having consumed it silently. Enrolment is idempotent + // only in the sense that the token is spent — a redelivery gets the refusal, which is + // correct and visible, where a lost request is neither. + _ = delivery.Ack(false) +} + +func (s *Server) reply(ctx context.Context, delivery amqp.Delivery, reply EnrolReply) { + if delivery.ReplyTo == "" { + s.log.Print("an enrolment request named no reply queue, so nothing can be told the answer") + return + } + body, err := json.Marshal(reply) + if err != nil { + s.log.Printf("cannot encode a reply: %v", err) + return + } + timeout, cancel := context.WithTimeout(ctx, 10*time.Second) + defer cancel() + + if err := s.channel.PublishWithContext(timeout, "", delivery.ReplyTo, false, false, + amqp.Publishing{ + ContentType: "application/json", + CorrelationId: delivery.CorrelationId, + Body: body, + }); err != nil { + s.log.Printf("cannot reply to %s: %v", delivery.ReplyTo, err) + } +}