A store row keeps its seat's protocol, and a holder may take work from its queue #108
@@ -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
@@ -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" }
|
||||||
} }
|
} }
|
||||||
|
|||||||
@@ -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)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -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}",
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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
|
||||||
|
|||||||
Reference in New Issue
Block a user