From c5dc7e732a58e2278e76fcdc47fb9c81dd31b210 Mon Sep 17 00:00:00 2001 From: jochen Date: Mon, 28 Sep 2026 02:08:56 +0200 Subject: [PATCH] A store row keeps its seat's protocol, and a holder may take work from its queue MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The seat table has name, scope, delivers and decision, and the protocol ADR 0129 gave a seat lives only in the compiled defaults; loading the rows dropped it, so no role's work queue was ever raised and the first build submitted over the new bus met "no response from stream". Until the table gains the columns, a row with no protocol keeps the compiled one of its name. And the holder of a seat is granted what taking work from its queue needs — asking about the worker consumer it binds, and acknowledging on it — which the first machine to try was refused. The control plane's own seat placeholders no longer include the old bus's port, which the switch removed with the variable. --- internal/broker/nats.go | 8 +++++++- internal/broker/testdata/composed.conf | 2 +- internal/catalogue/holdings_test.go | 15 +++++++++++++++ internal/catalogue/seat_into_test.go | 6 +----- internal/catalogue/seats.go | 25 +++++++++++++++++++++++-- 5 files changed, 47 insertions(+), 9 deletions(-) diff --git a/internal/broker/nats.go b/internal/broker/nats.go index 6ae882d..492a111 100644 --- a/internal/broker/nats.go +++ b/internal/broker/nats.go @@ -294,7 +294,13 @@ func PermissionsFor(p Principal) (Permissions, error) { // grants nothing anybody can use. sub = append(sub, "_DELIVER."+consumerDurable(p)) for _, s := range p.Holds { - sub = append(sub, "_DELIVER.SEAT_"+upperSnake(s.Name)+"_worker") + // Taking work from the role's queue: the worker consumer it binds (asked about, + // delivered on, acknowledged), each on the seat's own stream. The first machine to + // take work over the new bus was refused the asking (2026-09-28). + worker := "SEAT_" + upperSnake(s.Name) + "_worker" + stream := seatStreamName(s.Name) + sub = append(sub, "_DELIVER."+worker) + pub = append(pub, "$JS.API.CONSUMER.INFO."+stream+"."+worker, "$JS.ACK."+stream+"."+worker+".>") for _, a := range s.Accepts { sub = append(sub, seatSubject(s, "accept", a)) } diff --git a/internal/broker/testdata/composed.conf b/internal/broker/testdata/composed.conf index 9ebd7e8..365f227 100644 --- a/internal/broker/testdata/composed.conf +++ b/internal/broker/testdata/composed.conf @@ -37,7 +37,7 @@ accounts { subscribe: { allow: ["_DELIVER.one", "_INBOX.node.one.>", "mesh.node.one.declare"] } } } { user: "one.telegram", password: "$2a$11$tttttttttttttttttttttt", permissions: { - publish: { allow: ["$JS.ACK.EVENTS.one_telegram.>", "mesh.seat.telegram-sender.event.delivered", "mesh.seat.telegram-sender.event.failed"] } + publish: { allow: ["$JS.ACK.EVENTS.one_telegram.>", "$JS.ACK.SEAT_TELEGRAM_SENDER.SEAT_TELEGRAM_SENDER_worker.>", "$JS.API.CONSUMER.INFO.SEAT_TELEGRAM_SENDER.SEAT_TELEGRAM_SENDER_worker", "mesh.seat.telegram-sender.event.delivered", "mesh.seat.telegram-sender.event.failed"] } subscribe: { allow: ["_DELIVER.SEAT_TELEGRAM_SENDER_worker", "_DELIVER.one_telegram", "_INBOX.one.telegram.>", "mesh.mod.telegram.tool.status", "mesh.seat.telegram-sender.accept.send"] } allow_responses: { max: 1, ttl: "1m" } } } diff --git a/internal/catalogue/holdings_test.go b/internal/catalogue/holdings_test.go index 6394447..3c8d8eb 100644 --- a/internal/catalogue/holdings_test.go +++ b/internal/catalogue/holdings_test.go @@ -143,3 +143,18 @@ func TestAMembershipForTheNewBusIsComposedAsASealedFile(t *testing.T) { } } } + +// The store's seat rows have no protocol columns yet; loading them must not drop the protocol the +// bus is derived from, or no role's work queue is ever raised (found live, 2026-09-28). +func TestAStoreRowWithoutAProtocolKeepsTheCompiledOne(t *testing.T) { + was := Seats() + t.Cleanup(func() { UseSeats(was) }) + UseSeats([]Seat{{Name: "mesh-build-machine", Scope: ScopeMesh, Decision: "row"}}) + got, ok := SeatNamed("mesh-build-machine") + if !ok || len(got.Accepts) == 0 { + t.Fatalf("the build machine's seat lost what it accepts when loaded from the store: %+v", got) + } + if got.Decision != "row" { + t.Fatalf("the store's own columns were not kept: %+v", got) + } +} diff --git a/internal/catalogue/seat_into_test.go b/internal/catalogue/seat_into_test.go index bb7dc6b..04b911a 100644 --- a/internal/catalogue/seat_into_test.go +++ b/internal/catalogue/seat_into_test.go @@ -13,7 +13,6 @@ func TestASeatPlaceholderAnswersWhereThisMachinePutTheHolder(t *testing.T) { "type": "container", "id": "server", "name": "mesh-controller", "env": map[string]any{ "MESH_STORE_INVENTORY_PORT": "${seat:mesh-store:5432}", - "MESH_BROKER_AMQP_PORT": "${seat:mesh-broker:5672}", "MESH_BROKER_ADDRESS_PORT": "${seat:mesh-broker:5671}", "MESH_STORE_INVENTORY_FILE": "/run/secrets/inventory", }, @@ -27,7 +26,6 @@ func TestASeatPlaceholderAnswersWhereThisMachinePutTheHolder(t *testing.T) { env := control["env"].(map[string]any) for key, want := range map[string]string{ "MESH_STORE_INVENTORY_PORT": "6852", - "MESH_BROKER_AMQP_PORT": "5679", "MESH_BROKER_ADDRESS_PORT": "5671", "MESH_STORE_INVENTORY_FILE": "/run/secrets/inventory", } { @@ -146,7 +144,6 @@ func TestTheControlPlanesOwnAddressesFollowTheNodesPorts(t *testing.T) { "MESH_STORE_INVENTORY_PORT": "6852", "MESH_STORE_IDENTITY_PORT": "6852", "MESH_STORE_LICENCES_PORT": "6852", - "MESH_BROKER_AMQP_PORT": "5679", "MESH_BROKER_MANAGEMENT_PORT": "15673", "MESH_BROKER_ADDRESS_PORT": "5671", } { @@ -164,7 +161,7 @@ func TestTheControlPlanesOwnAddressesFollowTheNodesPorts(t *testing.T) { t.Fatal(err) } env, _ = fileNamed(out, "mesh-controller.server")["env"].(map[string]any) - if env["MESH_STORE_INVENTORY_PORT"] != "" || env["MESH_BROKER_AMQP_PORT"] != "" { + if env["MESH_STORE_INVENTORY_PORT"] != "" { t.Errorf("with no settings, the control plane is told %v", env) } } @@ -175,7 +172,6 @@ var SeatPorts = map[string]string{ "MESH_STORE_INVENTORY_PORT": "${seat:mesh-store:5432}", "MESH_STORE_IDENTITY_PORT": "${seat:mesh-store:5432}", "MESH_STORE_LICENCES_PORT": "${seat:mesh-store:5432}", - "MESH_BROKER_AMQP_PORT": "${seat:mesh-broker:5672}", "MESH_BROKER_MANAGEMENT_PORT": "${seat:mesh-broker:15672}", "MESH_BROKER_ADDRESS_PORT": "${seat:mesh-broker:5671}", } diff --git a/internal/catalogue/seats.go b/internal/catalogue/seats.go index 6cb9827..8d6c873 100644 --- a/internal/catalogue/seats.go +++ b/internal/catalogue/seats.go @@ -109,9 +109,30 @@ func DefaultSeats() []Seat { return append([]Seat(nil), defaultSeats...) } // than running on the set the binary shipped with. So the store can only ever *replace* the set with // a non-empty one, never erase it. func UseSeats(s []Seat) { - if len(s) > 0 { - seats = s + if len(s) == 0 { + return } + // **The store's rows carry no protocol yet, and the protocol is what the bus is derived + // from.** ADR 0129 gives a seat what it accepts, emits and serves; ADR 0122 moved the set into + // a table that has name, scope, delivers and decision and nothing else, and the columns for + // the rest are not there yet. So a row replacing a compiled entry would silently drop the + // protocol, and the roles' work queues would never be raised — found live as "no response + // from stream" the first time a build was submitted over the new bus (2026-09-28). Until the + // table gains the columns, a row without a protocol keeps the compiled one of the same name. + byName := map[string]Seat{} + for _, d := range defaultSeats { + byName[d.Name] = d + } + merged := make([]Seat, 0, len(s)) + for _, row := range s { + if len(row.Accepts)+len(row.Emits)+len(row.Serves) == 0 { + if d, known := byName[row.Name]; known { + row.Accepts, row.Emits, row.Serves = d.Accepts, d.Emits, d.Serves + } + } + merged = append(merged, row) + } + seats = merged } // aliases maps a seat's former names to its current canonical name (novox/hq ADR 0122). Loaded from