The outbound half went behind `Bus` and a node's two statements stopped naming a transport. This is the other half, and where the transport reached furthest: the run loop selected on a channel of the client library's own delivery type, so every part of holding a node in its mesh knew which bus it was on. `Link` is dialling, hearing and saying in one interface, because dialling is where the transport is chosen and choosing it twice is how one half of a node ends up on a different bus from the other. `Declaration` has one way of being done rather than two: a declaration set aside for a newer one is settled exactly as an applied one is, on both buses, and the difference is a fact the report carries. Four things this settled. **The host declares nothing on the new bus.** On the bus the mesh has it declares its own queue, because a queue that is not there means a node that hears nothing. Here it binds to a consumer the mesh made when the node enrolled, and a missing one is said as the mesh's to answer rather than quietly created with whatever this client happens to default to. **The pin is easier here than in the tool runtime, not harder.** The Go client takes a *tls.Config, so the same PinnedConfig with the same VerifyPeerCertificate does the work — the subject-alternative-name constraint recorded against the runtime's client is that client's, because it takes PEM strings with no verify hook. A host checks the fingerprint and nothing else. **Binding needs the subject as well as the consumer.** An empty subject is refused rather than taken to mean "whatever that consumer delivers", which the server said plainly and only when asked. **Reconnection stays the caller's.** Hold already decides when to try again and how long to wait; a client reconnecting underneath it would make that reasoning a duplicate of the library's. The drain keeps its live half and loses its catch-up half, as it said it would: verified that three declarations pushed to an absent node leave one on the stream, and it is the newest. One test-harness lesson worth the comment it got: delete-then-add is not a reset. A test that did that inherited the previous test's messages, and the symptom was a declaration counted as delivered twice — which reads as a redelivery bug in the code under test rather than as a dirty stream.
165 lines
4.9 KiB
Go
165 lines
4.9 KiB
Go
package link
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"fmt"
|
|
"net/url"
|
|
"time"
|
|
|
|
amqp "github.com/rabbitmq/amqp091-go"
|
|
)
|
|
|
|
// The host's link on the bus the mesh runs on today.
|
|
//
|
|
// Moved out of the run loop rather than changed: the dial, the queue, the prefetch window and the
|
|
// return handler are what they were, because the mesh is running on this and a bus nothing speaks
|
|
// yet is no reason to alter the one every node is on.
|
|
|
|
// currentLink is this node's connection as a channel.
|
|
type currentLink struct {
|
|
conn *amqp.Connection
|
|
channel *amqp.Channel
|
|
arrived chan Declaration
|
|
lost chan error
|
|
}
|
|
|
|
func dialCurrent(ctx context.Context, m Membership, timeout time.Duration) (Link, error) {
|
|
config, err := PinnedConfig(m.Fingerprint)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
dsn := fmt.Sprintf("amqps://%s:%s@%s/",
|
|
url.QueryEscape(m.Node), url.QueryEscape(m.Password), m.Broker)
|
|
conn, err := amqp.DialConfig(dsn, amqp.Config{
|
|
TLSClientConfig: config,
|
|
Dial: amqp.DefaultDial(timeout),
|
|
// Kept short so a node that has silently lost its route notices, rather than holding a
|
|
// connection the broker forgot about and believing it is still in the mesh.
|
|
Heartbeat: 10 * time.Second,
|
|
})
|
|
if err != nil {
|
|
if errors.Is(err, ErrWrongCertificate) {
|
|
return nil, err
|
|
}
|
|
return nil, fmt.Errorf("cannot reach the broker at %s: %w", m.Broker, err)
|
|
}
|
|
|
|
channel, err := conn.Channel()
|
|
if err != nil {
|
|
conn.Close()
|
|
return nil, err
|
|
}
|
|
|
|
queue := QueueFor(m.Node)
|
|
if _, err := channel.QueueDeclare(queue, true, false, false, false, nil); err != nil {
|
|
conn.Close()
|
|
return nil, fmt.Errorf("cannot declare this node's queue %s: %w", queue, err)
|
|
}
|
|
|
|
// Applying is one at a time — two at once would race on the same filesystem — but SEEING is
|
|
// not: with a prefetch of one the host could never know that a newer declaration was already
|
|
// waiting, and so applied every one of a backlog in turn, at the better part of a minute each,
|
|
// becoming things nobody wanted any more (novox/hq issue 031). A window of unacknowledged
|
|
// deliveries lets it drain to the newest; each declaration still survives a restart on the
|
|
// broker until it is acknowledged, which happens only after it is applied or set aside.
|
|
if err := channel.Qos(drainDepth, 0, false); err != nil {
|
|
conn.Close()
|
|
return nil, err
|
|
}
|
|
|
|
deliveries, err := channel.ConsumeWithContext(ctx, queue, "", false, false, false, false, nil)
|
|
if err != nil {
|
|
conn.Close()
|
|
return nil, err
|
|
}
|
|
|
|
l := ¤tLink{
|
|
conn: conn, channel: channel,
|
|
arrived: make(chan Declaration, drainDepth),
|
|
lost: make(chan error, 1),
|
|
}
|
|
|
|
// Published mandatory, so the broker hands back anything it cannot route rather than dropping
|
|
// it. Without this a report goes to an exchange with no matching binding, the publisher is told
|
|
// nothing, and the mesh believes this node never answered while the node believes it did —
|
|
// which is what happened when `report` was left unbound on the other side.
|
|
returned := channel.NotifyReturn(make(chan amqp.Return, 4))
|
|
go func() {
|
|
for r := range returned {
|
|
select {
|
|
case l.lost <- fmt.Errorf("the broker could not route this node's %s: %s (%d %s)",
|
|
r.RoutingKey, r.Exchange, r.ReplyCode, r.ReplyText):
|
|
default:
|
|
}
|
|
}
|
|
}()
|
|
|
|
closed := conn.NotifyClose(make(chan *amqp.Error, 1))
|
|
go func() {
|
|
select {
|
|
case reason := <-closed:
|
|
select {
|
|
case l.lost <- fmt.Errorf("the link closed: %v", reason):
|
|
default:
|
|
}
|
|
case <-ctx.Done():
|
|
}
|
|
}()
|
|
|
|
// One goroutine turning the library's deliveries into the mesh's words, so the run loop selects
|
|
// on one kind of thing whichever bus it is on.
|
|
go func() {
|
|
defer close(l.arrived)
|
|
for {
|
|
select {
|
|
case <-ctx.Done():
|
|
return
|
|
case delivery, ok := <-deliveries:
|
|
if !ok {
|
|
select {
|
|
case l.lost <- errors.New("the broker stopped delivering"):
|
|
default:
|
|
}
|
|
return
|
|
}
|
|
select {
|
|
case l.arrived <- currentDeclaration{delivery}:
|
|
case <-ctx.Done():
|
|
return
|
|
}
|
|
}
|
|
}
|
|
}()
|
|
|
|
return l, nil
|
|
}
|
|
|
|
func (l *currentLink) Declarations() <-chan Declaration { return l.arrived }
|
|
func (l *currentLink) Lost() <-chan error { return l.lost }
|
|
|
|
func (l *currentLink) Close() {
|
|
if l.channel != nil {
|
|
_ = l.channel.Close()
|
|
}
|
|
if l.conn != nil {
|
|
_ = l.conn.Close()
|
|
}
|
|
}
|
|
|
|
// Report and Alive are the outbound half, over the channel this link holds.
|
|
func (l *currentLink) Report(ctx context.Context, node string, body []byte) error {
|
|
return OverCurrent{Channel: l.channel}.Report(ctx, node, body)
|
|
}
|
|
|
|
func (l *currentLink) Alive(ctx context.Context, node string, body []byte) error {
|
|
return OverCurrent{Channel: l.channel}.Alive(ctx, node, body)
|
|
}
|
|
|
|
// currentDeclaration is one delivery from the bus the mesh has.
|
|
type currentDeclaration struct{ delivery amqp.Delivery }
|
|
|
|
func (d currentDeclaration) Body() []byte { return d.delivery.Body }
|
|
func (d currentDeclaration) Handled() error { return d.delivery.Ack(false) }
|