The controller follows the forge's merges (novox/hq 04-ISSUES/131). For each module recorded as built from that repository and branch it records the move to the merge commit and builds it — bases first, because a module built before the module it stands on is built against the old one and reports success, and a base that fails stops what stands on it. Nothing is pushed here: what a finished build does to the machines running the module stays the upgrade's decision. Two more things the same ordering gives: `build --behind` builds bases first, and `build --on <module>` rebuilds everything that stands on a module — the rebuild a changed base needs, which "behind" does not see because their sources did not move.
322 lines
11 KiB
Go
322 lines
11 KiB
Go
package link
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"fmt"
|
|
"time"
|
|
|
|
amqp "github.com/rabbitmq/amqp091-go"
|
|
)
|
|
|
|
// The consume side on the bus the mesh runs on today.
|
|
//
|
|
// Everything here was the serving loop's until the seam went in: the queues, the binds, the
|
|
// prefetch, and the list of messages the store could not take yet. It moved rather than changed —
|
|
// the behaviour this transport has is the behaviour it had, because the mesh is running on it and
|
|
// a bus nothing speaks yet is no reason to alter the one every node is on (ADR 0116).
|
|
|
|
// Prefetch is how many messages the bus hands the controller before it has settled them.
|
|
//
|
|
// More than one because a message the store could not take is held, unsettled, while the loop
|
|
// goes on answering others — an enrolment above all, which a host is waiting on (novox/hq issue
|
|
// 083). Bounded, because what is held is also what the bus has not kept on its own disk as
|
|
// pending.
|
|
const Prefetch = 64
|
|
|
|
// PrefetchHeadroom is how much of the prefetch is never held, so the loop always has messages to
|
|
// answer — an enrolment above all — while others wait for the store.
|
|
const PrefetchHeadroom = 8
|
|
|
|
// TryAgainAfter is how often the held are looked at. A store comes back in seconds, and a report a
|
|
// few seconds late is still current.
|
|
const TryAgainAfter = 2 * time.Second
|
|
|
|
// currentInbound consumes what nodes say over the bus the mesh has.
|
|
type currentInbound struct {
|
|
conn *amqp.Connection
|
|
channel *amqp.Channel
|
|
// upgrades and catchups are bound only when something is listening (Also).
|
|
upgrades bool
|
|
catchups bool
|
|
// held is every message the store could not take, by the bus's own delivery tag. Kept here
|
|
// rather than in the serving loop because holding a delivery unacknowledged is this
|
|
// transport's way of keeping it, and the other's is to hand it back to the server.
|
|
held map[uint64]*holding
|
|
// again is how often the held are looked at; zero means TryAgainAfter. Set by tests.
|
|
again time.Duration
|
|
}
|
|
|
|
// holding is one message kept for the store, and when to try it again.
|
|
type holding struct {
|
|
message *currentControl
|
|
due time.Time
|
|
about string
|
|
}
|
|
|
|
// Current is the consume side of the bus the mesh runs on today.
|
|
func Current(conn *amqp.Connection, channel *amqp.Channel) Inbound {
|
|
return ¤tInbound{conn: conn, channel: channel, held: map[uint64]*holding{}}
|
|
}
|
|
|
|
// Also binds the queue one more kind arrives on.
|
|
//
|
|
// The kinds nodes publish all share one queue and are bound at Connect, because a node may
|
|
// publish any of them and binding one while forgetting another is a message the bus accepts, finds
|
|
// no queue for, and drops — the publisher sees success and the consumer sees nothing. The two that
|
|
// are events get their own queue each, and only when something is listening.
|
|
func (c *currentInbound) Also(kind string) error {
|
|
switch kind {
|
|
case KindSourceMoved:
|
|
// Not followed on the bus the mesh is leaving: the forge's merges are announced on the
|
|
// new one, and this transport goes with the move (design 28, task 5.5).
|
|
return nil
|
|
case KindModuleMoved:
|
|
if err := c.bindEvent(UpgradeQueue, KeyModuleUpgraded); err != nil {
|
|
return err
|
|
}
|
|
c.upgrades = true
|
|
case KindCatchUp:
|
|
if err := c.bindEvent(CatchUpQueue, KeyCatchingUp); err != nil {
|
|
return err
|
|
}
|
|
c.catchups = true
|
|
default:
|
|
return fmt.Errorf("nothing binds a queue for %s on this bus", kind)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (c *currentInbound) bindEvent(queue, key string) error {
|
|
if _, err := c.channel.QueueDeclare(queue, true, false, false, false, nil); err != nil {
|
|
return fmt.Errorf("cannot declare the %s queue: %w", queue, err)
|
|
}
|
|
if err := c.channel.QueueBind(queue, key, EventsExchange, false, nil); err != nil {
|
|
return fmt.Errorf("cannot bind %s to %s/%s: %w", queue, EventsExchange, key, err)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (c *currentInbound) Close() {}
|
|
|
|
// Receive consumes until the context ends.
|
|
//
|
|
// One consumer per queue, deliberately: with two on one queue the bus would round-robin between
|
|
// them and each would receive half of what it expects — a fault this project has already had,
|
|
// between a module's daemon and its capability server.
|
|
func (c *currentInbound) Receive(ctx context.Context, act func(context.Context, Control)) error {
|
|
// A bounded prefetch rather than one. The loop still takes messages one at a time; what the
|
|
// prefetch buys is that a message the store could not take can be held while the loop goes on
|
|
// to the next, instead of every enrolment waiting behind it (novox/hq issue 083). Anything
|
|
// held goes back to the bus if the controller stops, because nothing held is acknowledged.
|
|
if err := c.channel.Qos(Prefetch, 0, false); err != nil {
|
|
return err
|
|
}
|
|
|
|
deliveries, err := c.channel.ConsumeWithContext(ctx, ControlQueue, "control-plane",
|
|
false, false, false, false, nil)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
// Its own queue and its own consumer for each event, for the reason above: two consumers on
|
|
// one queue split its messages between them, and an upgrade or a catch-up request going to
|
|
// whichever half was not listening is a gap that looks like a working mesh.
|
|
var upgrades, catchups <-chan amqp.Delivery
|
|
if c.upgrades {
|
|
upgrades, err = c.channel.ConsumeWithContext(ctx, UpgradeQueue, "control-plane-upgrades",
|
|
false, false, false, false, nil)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
}
|
|
if c.catchups {
|
|
catchups, err = c.channel.ConsumeWithContext(ctx, CatchUpQueue, "control-plane-catchup",
|
|
false, false, false, false, nil)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
}
|
|
|
|
closed := c.conn.NotifyClose(make(chan *amqp.Error, 1))
|
|
|
|
again := c.again
|
|
if again == 0 {
|
|
again = TryAgainAfter
|
|
}
|
|
ticker := time.NewTicker(again)
|
|
defer ticker.Stop()
|
|
|
|
for {
|
|
select {
|
|
case <-ctx.Done():
|
|
return nil
|
|
case <-ticker.C:
|
|
if ctx.Err() != nil {
|
|
return nil
|
|
}
|
|
c.retryHeld(ctx, act)
|
|
case delivery, ok := <-catchups:
|
|
if !ok {
|
|
if catchups != nil {
|
|
return errors.New("the bus stopped delivering catch-up requests")
|
|
}
|
|
continue
|
|
}
|
|
act(ctx, c.wrap(KindCatchUp, delivery))
|
|
case delivery, ok := <-upgrades:
|
|
// A nil channel blocks for ever, so this case simply never fires when nothing is
|
|
// listening for upgrades. Closed is different, and means the bus stopped.
|
|
if !ok {
|
|
if upgrades != nil {
|
|
return errors.New("the bus stopped delivering upgrades")
|
|
}
|
|
continue
|
|
}
|
|
act(ctx, c.wrap(KindModuleMoved, delivery))
|
|
case reason := <-closed:
|
|
// Said rather than returned quietly. A controller whose bus connection dropped is a
|
|
// mesh where nothing can be told anything, and the reason is the first thing anybody
|
|
// will want.
|
|
return fmt.Errorf("the bus connection closed: %v", reason)
|
|
case delivery, ok := <-deliveries:
|
|
if !ok {
|
|
return errors.New("the bus stopped delivering")
|
|
}
|
|
kind, known := kindOfKey[delivery.RoutingKey]
|
|
if !known {
|
|
// Rejected without requeue: a message nothing understands will not be understood
|
|
// on the next attempt either, and requeuing it would spin.
|
|
_ = delivery.Reject(false)
|
|
continue
|
|
}
|
|
act(ctx, c.wrap(kind, delivery))
|
|
}
|
|
}
|
|
}
|
|
|
|
// kindOfKey is how this transport's addressing becomes what the mesh calls a message.
|
|
var kindOfKey = map[string]string{
|
|
KeyEnrol: KindEnrolment,
|
|
KeyReport: KindReport,
|
|
KeyAlive: KindHeartbeat,
|
|
KeyBuilt: KindBuilt,
|
|
KeyModuleUpgraded: KindModuleMoved,
|
|
KeyCatchingUp: KindCatchUp,
|
|
}
|
|
|
|
func (c *currentInbound) wrap(kind string, delivery amqp.Delivery) *currentControl {
|
|
return ¤tControl{kind: kind, delivery: delivery, on: c}
|
|
}
|
|
|
|
// retryHeld hands every message whose delay has passed back to the loop. Each handler holds it
|
|
// again, settles it, or lets it go past the bound.
|
|
func (c *currentInbound) retryHeld(ctx context.Context, act func(context.Context, Control)) {
|
|
now := time.Now()
|
|
due := make([]*currentControl, 0, len(c.held))
|
|
for _, h := range c.held {
|
|
if !h.due.After(now) {
|
|
due = append(due, h.message)
|
|
}
|
|
}
|
|
for _, m := range due {
|
|
if ctx.Err() != nil {
|
|
return
|
|
}
|
|
act(ctx, m)
|
|
}
|
|
}
|
|
|
|
// currentControl is one delivery from the bus the mesh has, as the controller reads it.
|
|
type currentControl struct {
|
|
kind string
|
|
delivery amqp.Delivery
|
|
on *currentInbound
|
|
// about is what this message is about, as the handler named it; empty until it does.
|
|
about string
|
|
// first is when this message was first held for the store; zero while it has not been.
|
|
first time.Time
|
|
}
|
|
|
|
func (m *currentControl) Kind() string { return m.kind }
|
|
func (m *currentControl) Body() []byte { return m.delivery.Body }
|
|
func (m *currentControl) Redelivered() bool { return m.delivery.Redelivered }
|
|
|
|
func (m *currentControl) HeldFor() time.Duration {
|
|
if m.first.IsZero() {
|
|
return 0
|
|
}
|
|
return time.Since(m.first)
|
|
}
|
|
|
|
// Answer publishes to the reply queue the request named.
|
|
func (m *currentControl) Answer(ctx context.Context, body []byte) error {
|
|
if m.delivery.ReplyTo == "" {
|
|
return errors.New("that request named no reply queue, so nothing can be told the answer")
|
|
}
|
|
return m.on.channel.PublishWithContext(ctx, "", m.delivery.ReplyTo, false, false,
|
|
amqp.Publishing{
|
|
ContentType: "application/json",
|
|
CorrelationId: m.delivery.CorrelationId,
|
|
Body: body,
|
|
})
|
|
}
|
|
|
|
func (m *currentControl) Took() error {
|
|
m.forget()
|
|
return m.delivery.Ack(false)
|
|
}
|
|
|
|
// Drop rejects without requeue: on this bus that is what "understood, and not worth another
|
|
// attempt" is spelled as, and it is what feeds a dead-letter queue where one is configured.
|
|
func (m *currentControl) Drop() error {
|
|
m.forget()
|
|
return m.delivery.Reject(false)
|
|
}
|
|
|
|
// About names what this message is about, and lets go of whatever is held about the same thing:
|
|
// the held one is the past, and acting on it after this one would undo this one. Acknowledged
|
|
// rather than left to come back, because a held message nothing will act on is a place in the
|
|
// prefetch nothing gets back.
|
|
func (m *currentControl) About(what string) {
|
|
m.about = what
|
|
if what == "" {
|
|
return
|
|
}
|
|
for tag, h := range m.on.held {
|
|
if h.about != what || tag == m.delivery.DeliveryTag {
|
|
continue
|
|
}
|
|
delete(m.on.held, tag)
|
|
_ = h.message.delivery.Ack(false)
|
|
}
|
|
}
|
|
|
|
// Hold keeps the message unacknowledged and sets it aside to be handed back after the delay.
|
|
//
|
|
// Held no further than the prefetch leaves room: past that the bus would hand the loop nothing
|
|
// new — enrolments included — until something held was let go. A message that cannot be held says
|
|
// so, and the handler settles it its own way.
|
|
func (m *currentControl) Hold(after time.Duration) error {
|
|
if m.on.held == nil {
|
|
m.on.held = map[uint64]*holding{}
|
|
}
|
|
if _, already := m.on.held[m.delivery.DeliveryTag]; !already {
|
|
if len(m.on.held) >= Prefetch-PrefetchHeadroom {
|
|
return fmt.Errorf("%d messages are already held for the store, and holding more "+
|
|
"would stop the queue", len(m.on.held))
|
|
}
|
|
m.first = time.Now()
|
|
}
|
|
m.on.held[m.delivery.DeliveryTag] = &holding{
|
|
message: m, due: time.Now().Add(after), about: m.about,
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (m *currentControl) forget() {
|
|
if m.on != nil {
|
|
delete(m.on.held, m.delivery.DeliveryTag)
|
|
}
|
|
}
|