Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
da31bcb11e | ||
|
|
f03e7b33c9 |
@@ -753,8 +753,31 @@ func raiseTheBus(ctx context.Context, inv *inventory.Inventory, address string)
|
|||||||
if err := broker.RaiseSeats(js, inventory.MeshSeats(), holders); err != nil {
|
if err := broker.RaiseSeats(js, inventory.MeshSeats(), holders); err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
fmt.Printf("the bus at %s has its streams, and %d machine(s) can hear a declaration\n",
|
// And how every module hears what it consumes. Derived from the same records the user list is
|
||||||
broker.BareAddress(address), len(names))
|
// composed from, so a module the mesh grants a consumer's subjects has that consumer waiting.
|
||||||
|
// Done on every raise, not only when a credential is issued: every module moved onto this bus
|
||||||
|
// by the rollout was issued on the old one, and came up with nothing to bind to (2026-09-28).
|
||||||
|
records, err := inv.BusRecords(ctx)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
users, err := broker.Users(records)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
hearing := 0
|
||||||
|
for _, p := range users {
|
||||||
|
consumer, needed := broker.ConsumerFor(p)
|
||||||
|
if !needed {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
if err := js.EnsureConsumer(consumer); err != nil {
|
||||||
|
return fmt.Errorf("how %s on %s hears what it consumes: %w", p.Module, p.Node, err)
|
||||||
|
}
|
||||||
|
hearing++
|
||||||
|
}
|
||||||
|
fmt.Printf("the bus at %s has its streams, %d machine(s) can hear a declaration, and %d module(s) "+
|
||||||
|
"can hear what they consume\n", broker.BareAddress(address), len(names), hearing)
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
+11
-4
@@ -297,11 +297,18 @@ func PermissionsFor(p Principal) (Permissions, error) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// 2c. Its own consumer, which it **pulls**: the runtime asks for the next message and is
|
||||||
|
// answered on its own inbox, so what it needs is to ask about the consumer and to ask it
|
||||||
|
// for messages — its own consumer's name, and no other's. Pulled rather than pushed
|
||||||
|
// because that is the one shape a runtime's client binds without creating anything; the
|
||||||
|
// controller and the hosts are pushed to. Named here rather than through ConsumerFor,
|
||||||
|
// which asks for these permissions to build the consumer and would ask forever. A
|
||||||
|
// subject for a consumer that turns out not to exist grants nothing anybody can use.
|
||||||
|
pub = append(pub,
|
||||||
|
"$JS.API.CONSUMER.INFO."+consumerStream(p)+"."+consumerDurable(p),
|
||||||
|
"$JS.API.CONSUMER.MSG.NEXT."+consumerStream(p)+"."+consumerDurable(p))
|
||||||
|
|
||||||
// 3. Seats it holds: full participation.
|
// 3. Seats it holds: full participation.
|
||||||
// Its consumer's name, not ConsumerFor: that asks for these permissions to build the
|
|
||||||
// consumer, and would ask forever. A subject for a consumer that turns out not to exist
|
|
||||||
// grants nothing anybody can use.
|
|
||||||
sub = append(sub, "_DELIVER."+consumerDurable(p))
|
|
||||||
for _, s := range p.Holds {
|
for _, s := range p.Holds {
|
||||||
// Taking work from the role's queue: the worker consumer it binds (asked about,
|
// 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
|
// delivered on, acknowledged), each on the seat's own stream. The first machine to
|
||||||
|
|||||||
@@ -338,3 +338,24 @@ func admits(pattern, subject []string) bool {
|
|||||||
}
|
}
|
||||||
return len(pattern) == len(subject)
|
return len(pattern) == len(subject)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// A module pulls its own consumer — asks about it, asks it for messages — and no other module's.
|
||||||
|
func TestAModulePullsItsOwnConsumerAndNoOthers(t *testing.T) {
|
||||||
|
p, _ := PermissionsFor(Principal{Kind: KindModule, Node: "one", Module: "audit",
|
||||||
|
Consumes: []string{"shop.order.placed"}, PasswordHash: "x"})
|
||||||
|
for _, want := range []string{"$JS.API.CONSUMER.INFO.EVENTS.one_audit", "$JS.API.CONSUMER.MSG.NEXT.EVENTS.one_audit"} {
|
||||||
|
if !slices.Contains(p.Publish, want) {
|
||||||
|
t.Errorf("a module cannot bind its own consumer: %v lacks %s", p.Publish, want)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
for _, s := range p.Publish {
|
||||||
|
if strings.Contains(s, "CONSUMER.") && !strings.HasSuffix(s, ".one_audit") {
|
||||||
|
t.Errorf("a module may reach another consumer: %s", s)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
for _, s := range p.Subscribe {
|
||||||
|
if strings.HasPrefix(s, "_DELIVER.") {
|
||||||
|
t.Errorf("a module is granted a push delivery it never binds: %s", s)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
+6
-6
@@ -37,18 +37,18 @@ 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.>", "$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"] }
|
publish: { allow: ["$JS.ACK.EVENTS.one_telegram.>", "$JS.ACK.SEAT_TELEGRAM_SENDER.SEAT_TELEGRAM_SENDER_worker.>", "$JS.API.CONSUMER.INFO.EVENTS.one_telegram", "$JS.API.CONSUMER.INFO.SEAT_TELEGRAM_SENDER.SEAT_TELEGRAM_SENDER_worker", "$JS.API.CONSUMER.MSG.NEXT.EVENTS.one_telegram", "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.>", "mesh.seat.telegram-sender.accept.send"] }
|
subscribe: { allow: ["_DELIVER.SEAT_TELEGRAM_SENDER_worker", "_INBOX.one.telegram.>", "mesh.mod.telegram.tool.>", "mesh.seat.telegram-sender.accept.send"] }
|
||||||
allow_responses: { max: 1, ttl: "1m" }
|
allow_responses: { max: 1, ttl: "1m" }
|
||||||
} }
|
} }
|
||||||
{ user: "two.audit", password: "$2a$11$aaaaaaaaaaaaaaaaaaaaaa", permissions: {
|
{ user: "two.audit", password: "$2a$11$aaaaaaaaaaaaaaaaaaaaaa", permissions: {
|
||||||
publish: { allow: ["$JS.ACK.EVENTS.two_audit.>"] }
|
publish: { allow: ["$JS.ACK.EVENTS.two_audit.>", "$JS.API.CONSUMER.INFO.EVENTS.two_audit", "$JS.API.CONSUMER.MSG.NEXT.EVENTS.two_audit"] }
|
||||||
subscribe: { allow: ["_DELIVER.two_audit", "_INBOX.two.audit.>", "mesh.mod.audit.tool.>", "mesh.mod.shop.event.order.placed"] }
|
subscribe: { allow: ["_INBOX.two.audit.>", "mesh.mod.audit.tool.>", "mesh.mod.shop.event.order.placed"] }
|
||||||
allow_responses: { max: 1, ttl: "1m" }
|
allow_responses: { max: 1, ttl: "1m" }
|
||||||
} }
|
} }
|
||||||
{ user: "two.shop", password: "$2a$11$ssssssssssssssssssssss", permissions: {
|
{ user: "two.shop", password: "$2a$11$ssssssssssssssssssssss", permissions: {
|
||||||
publish: { allow: ["$JS.ACK.EVENTS.two_shop.>", "mesh.mod.shop.event.order.placed", "mesh.seat.telegram-sender.accept.send"] }
|
publish: { allow: ["$JS.ACK.EVENTS.two_shop.>", "$JS.API.CONSUMER.INFO.EVENTS.two_shop", "$JS.API.CONSUMER.MSG.NEXT.EVENTS.two_shop", "mesh.mod.shop.event.order.placed", "mesh.seat.telegram-sender.accept.send"] }
|
||||||
subscribe: { allow: ["_DELIVER.two_shop", "_INBOX.two.shop.>", "mesh.mod.shop.tool.>"] }
|
subscribe: { allow: ["_INBOX.two.shop.>", "mesh.mod.shop.tool.>"] }
|
||||||
allow_responses: { max: 1, ttl: "1m" }
|
allow_responses: { max: 1, ttl: "1m" }
|
||||||
} }
|
} }
|
||||||
]
|
]
|
||||||
|
|||||||
Reference in New Issue
Block a user