Merge pull request 'A store row keeps its seat's protocol, and a holder may take work from its queue' (#108) from fix/store-seats-keep-their-protocol into main

This commit is contained in:
2026-09-28 02:21:24 +02:00
committed by jschoubben
6 changed files with 110 additions and 10 deletions
+63 -1
View File
@@ -5,6 +5,7 @@ import (
"encoding/json" "encoding/json"
"errors" "errors"
"fmt" "fmt"
"os"
"strings" "strings"
"time" "time"
@@ -33,7 +34,7 @@ import (
// ability to change things, not the services its modules are serving — measured on 2026-09-27, when // ability to change things, not the services its modules are serving — measured on 2026-09-27, when
// a seat emptied mid-change and the control plane looped for two hours while every service stayed up. // a seat emptied mid-change and the control plane looped for two hours while every service stayed up.
const rolloutUsage = "rollout check | rollout mint [--again] | rollout --confirm" const rolloutUsage = "rollout check | rollout mint [--again] | rollout hand <node> | rollout --confirm"
func rolloutCommand(ctx context.Context, args []string) error { func rolloutCommand(ctx context.Context, args []string) error {
switch { switch {
@@ -41,6 +42,8 @@ func rolloutCommand(ctx context.Context, args []string) error {
return rolloutCheck(ctx) return rolloutCheck(ctx)
case len(args) == 1 && args[0] == "mint": case len(args) == 1 && args[0] == "mint":
return rolloutMint(ctx, false) return rolloutMint(ctx, false)
case len(args) == 2 && args[0] == "hand":
return rolloutHand(ctx, args[1])
case len(args) == 2 && args[0] == "mint" && args[1] == "--again": case len(args) == 2 && args[0] == "mint" && args[1] == "--again":
// Every credential minted afresh, whether or not one exists — for a mint that was wrong // Every credential minted afresh, whether or not one exists — for a mint that was wrong
// before anything was pushed. Afterwards nothing that received the old one still works, // before anything was pushed. Afterwards nothing that received the old one still works,
@@ -387,3 +390,62 @@ func providesBus(m catalogue.Manifest) bool {
} }
return false return false
} }
// rolloutHand mints a machine its credential for the new bus afresh and prints its membership
// once, for an operator to carry by hand — the rescue for a machine that cannot be reached over
// any bus: rotated while it still held the old password, or reachable only by ssh. The plaintext
// exists on this terminal and then only where it is written; the store keeps the hash, and the
// sealed copy in the machine's declaration is replaced too, so the next push says the same.
func rolloutHand(ctx context.Context, node string) error {
open, err := openStores(ctx)
if err != nil {
return err
}
defer open.Close()
inv := open.inventory
known, err := broker.FromEnvironment()
if err != nil {
return fmt.Errorf("the bus's certificate is not known to this process: %w", err)
}
busAddress, _, err := broker.OnNATS()
if err != nil {
return err
}
if busAddress == "" {
return errors.New("this control plane is not on the new bus, so there is no membership to hand out")
}
_, _, bare := broker.CredentialIn(busAddress)
if _, after, has := strings.Cut(bare, "://"); has {
bare = after
}
if _, err := inv.NodeByName(ctx, node); err != nil {
return err
}
p := broker.Principal{Kind: broker.KindNode, Node: node}
password, err := inv.MintBusPassword(ctx, inventory.BusUser{Username: p.Username(), Kind: inventory.BusNode, Node: node})
if err != nil {
return err
}
membership, _ := json.Marshal(map[string]string{
"broker": bare, "fingerprint": known.Fingerprint, "password": password, "transport": "nats",
})
key, err := inv.SealingKeyOf(ctx, node)
if err != nil {
return err
}
sealed, err := secrets.Seal(key, membership)
if err != nil {
return err
}
if err := inv.PutBusMembership(ctx, node, sealed); err != nil {
return err
}
// The one line of output is the membership itself, so it can be piped to the machine without
// being read on the way. Everything else goes to stderr.
fmt.Fprintf(os.Stderr, "%s's credential is minted afresh. Write this to %s on it and restart its host; "+
"then push the machine running the bus so the user list carries the new hash.\n",
node, catalogue.BusMembershipPath)
fmt.Println(string(membership))
return nil
}
+7 -1
View File
@@ -294,7 +294,13 @@ func PermissionsFor(p Principal) (Permissions, error) {
// grants nothing anybody can use. // grants nothing anybody can use.
sub = append(sub, "_DELIVER."+consumerDurable(p)) sub = append(sub, "_DELIVER."+consumerDurable(p))
for _, s := range p.Holds { 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 { for _, a := range s.Accepts {
sub = append(sub, seatSubject(s, "accept", a)) sub = append(sub, seatSubject(s, "accept", a))
} }
+1 -1
View File
@@ -37,7 +37,7 @@ accounts {
subscribe: { allow: ["_DELIVER.one", "_INBOX.node.one.>", "mesh.node.one.declare"] } subscribe: { allow: ["_DELIVER.one", "_INBOX.node.one.>", "mesh.node.one.declare"] }
} } } }
{ user: "one.telegram", password: "$2a$11$tttttttttttttttttttttt", permissions: { { 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"] } 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" } allow_responses: { max: 1, ttl: "1m" }
} } } }
+15
View File
@@ -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)
}
}
+1 -5
View File
@@ -13,7 +13,6 @@ func TestASeatPlaceholderAnswersWhereThisMachinePutTheHolder(t *testing.T) {
"type": "container", "id": "server", "name": "mesh-controller", "type": "container", "id": "server", "name": "mesh-controller",
"env": map[string]any{ "env": map[string]any{
"MESH_STORE_INVENTORY_PORT": "${seat:mesh-store:5432}", "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_BROKER_ADDRESS_PORT": "${seat:mesh-broker:5671}",
"MESH_STORE_INVENTORY_FILE": "/run/secrets/inventory", "MESH_STORE_INVENTORY_FILE": "/run/secrets/inventory",
}, },
@@ -27,7 +26,6 @@ func TestASeatPlaceholderAnswersWhereThisMachinePutTheHolder(t *testing.T) {
env := control["env"].(map[string]any) env := control["env"].(map[string]any)
for key, want := range map[string]string{ for key, want := range map[string]string{
"MESH_STORE_INVENTORY_PORT": "6852", "MESH_STORE_INVENTORY_PORT": "6852",
"MESH_BROKER_AMQP_PORT": "5679",
"MESH_BROKER_ADDRESS_PORT": "5671", "MESH_BROKER_ADDRESS_PORT": "5671",
"MESH_STORE_INVENTORY_FILE": "/run/secrets/inventory", "MESH_STORE_INVENTORY_FILE": "/run/secrets/inventory",
} { } {
@@ -146,7 +144,6 @@ func TestTheControlPlanesOwnAddressesFollowTheNodesPorts(t *testing.T) {
"MESH_STORE_INVENTORY_PORT": "6852", "MESH_STORE_INVENTORY_PORT": "6852",
"MESH_STORE_IDENTITY_PORT": "6852", "MESH_STORE_IDENTITY_PORT": "6852",
"MESH_STORE_LICENCES_PORT": "6852", "MESH_STORE_LICENCES_PORT": "6852",
"MESH_BROKER_AMQP_PORT": "5679",
"MESH_BROKER_MANAGEMENT_PORT": "15673", "MESH_BROKER_MANAGEMENT_PORT": "15673",
"MESH_BROKER_ADDRESS_PORT": "5671", "MESH_BROKER_ADDRESS_PORT": "5671",
} { } {
@@ -164,7 +161,7 @@ func TestTheControlPlanesOwnAddressesFollowTheNodesPorts(t *testing.T) {
t.Fatal(err) t.Fatal(err)
} }
env, _ = fileNamed(out, "mesh-controller.server")["env"].(map[string]any) 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) 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_INVENTORY_PORT": "${seat:mesh-store:5432}",
"MESH_STORE_IDENTITY_PORT": "${seat:mesh-store:5432}", "MESH_STORE_IDENTITY_PORT": "${seat:mesh-store:5432}",
"MESH_STORE_LICENCES_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_MANAGEMENT_PORT": "${seat:mesh-broker:15672}",
"MESH_BROKER_ADDRESS_PORT": "${seat:mesh-broker:5671}", "MESH_BROKER_ADDRESS_PORT": "${seat:mesh-broker:5671}",
} }
+23 -2
View File
@@ -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 // 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. // a non-empty one, never erase it.
func UseSeats(s []Seat) { func UseSeats(s []Seat) {
if len(s) > 0 { if len(s) == 0 {
seats = s 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 // aliases maps a seat's former names to its current canonical name (novox/hq ADR 0122). Loaded from