Files
mesh-host/internal/link/asking_current.go
T
jschoubben 9072f60a30 Enrolment behind a seam, with both transports
The last of the host's link that still named a transport. `Asking` is one
enrolment conversation — a connection made with the token, a question asked, and
an answer waited for — and it is its own seam rather than part of `Link` because
almost nothing about it is the same: the credential is a one-time secret, there
is no declaration to hear, and a node that fails here is not in the mesh at all,
where a node that fails in `Link` has merely lost touch with one it belongs to.

`Enrol`'s thirteen arguments became an `Approach` — where, which certificate,
which bus — and the request it already had. The token says nothing about which
bus, and does not need to: every token names the one the mesh runs on today until
the rollout.

**The reply address is the whole of what changes on the new bus**, and it is
forced rather than preferred. Verified against a running server, both halves: the
answer reaches the node at the address its request carried in the payload, and
the transport's own reply field held something else entirely by the time the
consumer saw it — the consumer's ack address, exactly as design 25 §2 says. The
test asserts the field is *not* the node's inbox, so a future server that stopped
claiming it would fail this rather than let the reason quietly become folklore.

The inbox is under `_INBOX.enrol.<node>.`, which is exactly what the enrolling
user may subscribe and no wider, with a random tail per attempt: a reply left
over from an attempt that timed out is not the answer to this question, which is
what the correlation id does on the other transport. Subscribed before anything
is published, because a node that published first could miss an answer to a
question nobody was listening for.
2026-09-27 01:31:46 +02:00

130 lines
3.8 KiB
Go

package link
import (
"context"
"errors"
"fmt"
"net/url"
"time"
amqp "github.com/rabbitmq/amqp091-go"
)
// The enrolment conversation on the bus the mesh runs on today.
//
// Moved out of Enrol rather than changed. The reply travels on this node's own queue and is picked
// out by the correlation id the request carried, which is what this transport's reply field means
// and has always meant.
type currentAsking struct {
conn *amqp.Connection
channel *amqp.Channel
queue string
node string
replies <-chan amqp.Delivery
closed chan *amqp.Error
}
func presentCurrent(_ context.Context, to Approach, node, secret string,
timeout time.Duration) (Asking, error) {
config, err := PinnedConfig(to.Fingerprint)
if err != nil {
return nil, err
}
// The account name is the node's, and the password is the token's secret. Escaped because a name
// or secret containing a colon or an at-sign would otherwise change which host this connects to
// — a credential silently redirecting a connection is the worst shape this could take.
dsn := fmt.Sprintf("amqps://%s:%s@%s/",
url.QueryEscape(node), url.QueryEscape(secret), to.Address)
conn, err := amqp.DialConfig(dsn, amqp.Config{
TLSClientConfig: config,
Dial: amqp.DefaultDial(timeout),
})
if err != nil {
if errors.Is(err, ErrWrongCertificate) {
return nil, err
}
// Not quoted back: the DSN carries the one-time secret.
return nil, fmt.Errorf("cannot reach the broker at %s as %s: %w", to.Address, node, err)
}
channel, err := conn.Channel()
if err != nil {
conn.Close()
return nil, err
}
// This node's own queue, which its account is scoped to and nothing else may read.
queue, err := channel.QueueDeclare(QueueFor(node), true, false, false, false, nil)
if err != nil {
conn.Close()
return nil, fmt.Errorf("cannot declare this node's queue %s: %w", QueueFor(node), err)
}
replies, err := channel.Consume(queue.Name, "", true, false, false, false, nil)
if err != nil {
conn.Close()
return nil, err
}
return &currentAsking{
conn: conn, channel: channel, queue: queue.Name, node: node, replies: replies,
closed: conn.NotifyClose(make(chan *amqp.Error, 1)),
}, nil
}
func (a *currentAsking) Close() {
if a.channel != nil {
_ = a.channel.Close()
}
if a.conn != nil {
_ = a.conn.Close()
}
}
func (a *currentAsking) Ask(ctx context.Context, request []byte, wait time.Duration) ([]byte, error) {
correlation := fmt.Sprintf("%s-%d", a.node, time.Now().UnixNano())
publish, cancel := context.WithTimeout(ctx, wait)
defer cancel()
if err := a.channel.PublishWithContext(publish, Exchange, KeyEnrol, false, false,
amqp.Publishing{
ContentType: "application/json",
CorrelationId: correlation,
ReplyTo: a.queue,
Body: request,
}); err != nil {
return nil, fmt.Errorf("cannot publish to the %s exchange: %w", Exchange, err)
}
// Waited for rather than assumed. A published message that nothing answers means the control
// plane is not running, and a node that carried on regardless would believe it had joined a mesh
// that has never heard of it.
deadline := time.NewTimer(wait)
defer deadline.Stop()
for {
select {
case <-ctx.Done():
return nil, ctx.Err()
case reason := <-a.closed:
return nil, fmt.Errorf("the broker closed the connection: %v", reason)
case <-deadline.C:
return nil, fmt.Errorf(
"the broker accepted this node's connection and nothing answered within %s. The "+
"mesh's broker is running and its control plane is not", wait)
case delivery, ok := <-a.replies:
if !ok {
return nil, errors.New("the broker stopped delivering")
}
// Anything else on this queue is not the answer to this question.
if delivery.CorrelationId != correlation {
continue
}
return delivery.Body, nil
}
}
}