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) }