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