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 ¤tAsking{ 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 } } }