Step 3.4, first half. Every one of these took an *amqp.Channel, so the transport reached every caller and swapping it meant touching all of them. The seam turned out to be small — the controller sends exactly two kinds of message that expect no answer — which is the same measurement that said this bus could be replaced at all. Bus is stated in the mesh's words, not a transport's: PublishEvent and PublishDeclaration. Two implementations, both shipping, because steps 1 to 4 leave every node on AMQP and the NATS one is selected at the rollout. Both ship is also what makes them comparable: one conformance fixture holds both to the same envelope, and the NATS one is checked against a real server reading back from the stream rather than from the code that wrote it. Still on *amqp.Channel: RequestBuild and Ask, which carry reply-queue machinery, and the whole consume side — the control loop, enrolment, serve.
44 lines
1.2 KiB
Go
44 lines
1.2 KiB
Go
package link
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"fmt"
|
|
"time"
|
|
)
|
|
|
|
// Signer is whatever holds the control plane's signing key.
|
|
type Signer interface {
|
|
Sign(ctx context.Context, message []byte) ([]byte, error)
|
|
}
|
|
|
|
// Declare sends a node what it should be, signed.
|
|
//
|
|
// The signature is over the declaration exactly as it is published — the same bytes the node
|
|
// verifies. Anything that re-encoded between here and there would produce a signature over
|
|
// something else, and the node would refuse a declaration that was genuinely the mesh's.
|
|
//
|
|
// Published to the node's own queue, which its account alone may read.
|
|
func Declare(ctx context.Context, bus Bus, signer Signer, node string,
|
|
declaration []byte, timeout time.Duration) error {
|
|
|
|
if !json.Valid(declaration) {
|
|
return fmt.Errorf("refusing to send %s something that is not a declaration", node)
|
|
}
|
|
|
|
signature, err := signer.Sign(ctx, declaration)
|
|
if err != nil {
|
|
return fmt.Errorf("cannot sign a declaration for %s: %w", node, err)
|
|
}
|
|
|
|
body, err := json.Marshal(Signed{Declaration: declaration, Signature: signature})
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
publish, cancel := context.WithTimeout(ctx, timeout)
|
|
defer cancel()
|
|
|
|
return bus.PublishDeclaration(publish, node, body)
|
|
}
|