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