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