ADR 0121 carried through to working code. `Builders` is the asking side and `BuildMachine` the taking side, each with an implementation per bus, and the builder binary and the `build` command now go through them. On the bus being built, one publish does what two did. The old bus answered the asker through a reply queue and announced to an events exchange, because two audiences meant two topologies. Here the outcome is the role's own event: the asker matches it by the id its request carried, the controller records it, the catalogue places it in the graph. So a build machine publishes once, needs a reply queue for nothing, and needs a grant over nobody's inbox — which is what ruled out the alternatives. The outcome carries the module name now. Only the manifest says what was built, and on the old bus the separate announcement carried it; with one message for three readers it belongs in the result. A failed build names none, because it produced no module version and the catalogue would otherwise place something that was never made. Checked against a real server: the whole round trip; a third party on the role's event hearing the same outcome the asker did, which is the claim the decision rests on; work leaving the queue once settled, so no second machine repeats it; work submitted with no machine holding the role waiting instead of failing, and being done when one arrives; and work a machine handed back coming round again. One thing I got wrong twice now and have written down where it bit: binding to a consumer must name that consumer's own filter subject, not the narrower subject the caller cares about. The client compares the two and refuses anything that is not equal, with "subject does not match consumer".
185 lines
6.2 KiB
Go
185 lines
6.2 KiB
Go
package link
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"time"
|
|
|
|
amqp "github.com/rabbitmq/amqp091-go"
|
|
)
|
|
|
|
// The build flow on the bus the mesh runs on today.
|
|
//
|
|
// Moved behind the seam rather than changed. The queue, the reply binding and the correlation are
|
|
// what they were, because the mesh is running on this.
|
|
|
|
// currentBuilds asks for builds over a channel.
|
|
type currentBuilds struct{ channel *amqp.Channel }
|
|
|
|
// BuildsOverCurrent is the asking side on the bus the mesh has.
|
|
func BuildsOverCurrent(channel *amqp.Channel) Builders { return currentBuilds{channel: channel} }
|
|
|
|
func (b currentBuilds) Close() {}
|
|
|
|
func (b currentBuilds) Submit(ctx context.Context, request BuildRequest,
|
|
wait time.Duration) (BuildResult, error) {
|
|
|
|
// Its own queue for the answer, declared before the ask. Consuming from the shared exchange
|
|
// would mean competing with the controller's own consumer for a message meant for this caller.
|
|
replies, err := b.channel.QueueDeclare(ReplyQueue(request.ID), false, true, true, false, nil)
|
|
if err != nil {
|
|
return BuildResult{}, err
|
|
}
|
|
// **A builder never publishes to the default exchange**, because permission there is per
|
|
// exchange and not per queue — a builder allowed to use it could publish into any node's queue,
|
|
// which is the privilege a build machine most obviously should not have. The cost is that every
|
|
// asker sees every result, which is why the correlation is checked below rather than assumed.
|
|
if err := b.channel.QueueBind(replies.Name, KeyBuilt, Exchange, false, nil); err != nil {
|
|
return BuildResult{}, err
|
|
}
|
|
answers, err := b.channel.ConsumeWithContext(ctx, replies.Name, "", true, true, false, false, nil)
|
|
if err != nil {
|
|
return BuildResult{}, err
|
|
}
|
|
|
|
body, err := json.Marshal(request)
|
|
if err != nil {
|
|
return BuildResult{}, err
|
|
}
|
|
if err := b.channel.PublishWithContext(ctx, "", BuildQueue, false, false, amqp.Publishing{
|
|
ContentType: "application/json",
|
|
DeliveryMode: amqp.Persistent,
|
|
CorrelationId: request.ID,
|
|
ReplyTo: replies.Name,
|
|
Body: body,
|
|
}); err != nil {
|
|
return BuildResult{}, err
|
|
}
|
|
|
|
waiting, cancel := context.WithTimeout(ctx, wait)
|
|
defer cancel()
|
|
for {
|
|
select {
|
|
case <-waiting.Done():
|
|
return BuildResult{}, waitingFor(wait)
|
|
case delivery, ok := <-answers:
|
|
if !ok {
|
|
return BuildResult{}, errors.New("the connection closed while waiting for a build")
|
|
}
|
|
result, mine, err := theOutcomeOf(delivery.Body, request.ID)
|
|
if err != nil {
|
|
return BuildResult{}, err
|
|
}
|
|
if mine {
|
|
return result, nil
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
// --- the machine's side ---------------------------------------------------------------------
|
|
|
|
type currentMachine struct {
|
|
conn *amqp.Connection
|
|
channel *amqp.Channel
|
|
on string
|
|
}
|
|
|
|
// MachineOverCurrent takes build work over a channel.
|
|
func MachineOverCurrent(conn *amqp.Connection, channel *amqp.Channel, on string) BuildMachine {
|
|
return ¤tMachine{conn: conn, channel: channel, on: on}
|
|
}
|
|
|
|
func (m *currentMachine) Close() {}
|
|
|
|
func (m *currentMachine) Take(ctx context.Context, do func(context.Context, Build)) error {
|
|
if _, err := m.channel.QueueDeclare(BuildQueue, true, false, false, false, nil); err != nil {
|
|
return err
|
|
}
|
|
// One at a time. A machine that took five requests at once would run five container builds
|
|
// against one runtime and finish all of them slower than it would have finished the first — and
|
|
// the queue is what shares work between machines, so nothing is lost by it.
|
|
if err := m.channel.Qos(1, 0, false); err != nil {
|
|
return err
|
|
}
|
|
// Not auto-acknowledged: a request acknowledged on arrival is a build that vanishes if this
|
|
// process dies mid-way, with nobody waiting on it ever hearing why.
|
|
requests, err := m.channel.ConsumeWithContext(ctx, BuildQueue, "mesh-builder",
|
|
false, false, false, false, nil)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
for {
|
|
select {
|
|
case <-ctx.Done():
|
|
return nil
|
|
case delivery, ok := <-requests:
|
|
if !ok {
|
|
return errors.New("the broker closed the connection")
|
|
}
|
|
var request BuildRequest
|
|
if err := json.Unmarshal(delivery.Body, &request); err != nil {
|
|
// Unreadable: rejected rather than retried, because the next attempt reads the same
|
|
// bytes. Nobody waiting hears an answer, which is correct — there was no request.
|
|
_ = delivery.Reject(false)
|
|
continue
|
|
}
|
|
do(ctx, ¤tBuild{request: request, delivery: delivery, on: m.on, channel: m.channel})
|
|
}
|
|
}
|
|
}
|
|
|
|
type currentBuild struct {
|
|
request BuildRequest
|
|
delivery amqp.Delivery
|
|
on string
|
|
channel *amqp.Channel
|
|
}
|
|
|
|
func (b *currentBuild) Request() BuildRequest { return b.request }
|
|
|
|
// Announce answers and announces, which on this bus are two publishes to two exchanges.
|
|
//
|
|
// The reply goes to whoever asked, correlated to their request; the announcement says to the whole
|
|
// mesh that a module now exists at a commit (novox/hq ADR 0072). Only a successful build is
|
|
// announced: a failed one produced no module version, and announcing one would put something in the
|
|
// graph that was never made.
|
|
func (b *currentBuild) Announce(ctx context.Context, result BuildResult) error {
|
|
body, err := json.Marshal(result)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if err := b.channel.PublishWithContext(ctx, Exchange, KeyBuilt, false, false, amqp.Publishing{
|
|
ContentType: "application/json",
|
|
CorrelationId: result.ID,
|
|
Body: body,
|
|
}); err != nil {
|
|
return fmt.Errorf("cannot answer a build request: %w", err)
|
|
}
|
|
if result.Failed != "" || result.Commit == "" {
|
|
return nil
|
|
}
|
|
return EmitEvent(ctx, OverCurrent{Channel: b.channel}, KeyModuleBuilt, "builder", b.on,
|
|
announcementOf(result))
|
|
}
|
|
|
|
func (b *currentBuild) Done() error { return b.delivery.Ack(false) }
|
|
|
|
func (b *currentBuild) Hold(time.Duration) error {
|
|
// No delayed redelivery on this bus: handed back at once, which is what it has always done.
|
|
return b.delivery.Nack(false, true)
|
|
}
|
|
|
|
// announcementOf is what the mesh is told about a finished build. One function, so the two
|
|
// transports cannot describe the same build differently.
|
|
func announcementOf(result BuildResult) map[string]any {
|
|
return map[string]any{
|
|
"module": ModuleOf(result.Manifest), "commit": result.Commit,
|
|
"repository": result.Repository, "path": result.Path, "ref": result.Ref,
|
|
"manifest": json.RawMessage(result.Manifest), "against": result.Against,
|
|
"made": result.Made,
|
|
}
|
|
}
|