Compare commits

...
Author SHA1 Message Date
jschoubben 3756bb3460 A base the registry already holds is not pulled from upstream again
A base is named by digest, and a digest the mesh's registry holds under the module's repository
is the same bytes whatever upstream would say. Asked on every build, the public hub's anonymous
pull limit was reached on the first merge that rebuilt a whole catalogue, and every module whose
base lives there failed on a copy it did not need.
2026-09-28 04:48:32 +02:00
mesh-admin 60be9c5360 Merge pull request 'A module hears what it consumes: its consumer is raised with the bus, and it pulls it' (#118) from fix/a-module-hears-what-it-consumes into main 2026-09-28 02:29:34 +00:00
jschoubben da31bcb11e A module hears what it consumes: its consumer is raised with the bus, and it pulls it
Every module moved onto the bus by the rollout was issued on the old one, so none had a consumer
waiting; and the grant named a push delivery a runtime's client never binds, while the pull it
does make — asking about its consumer, asking it for messages — was refused. The consumers a
module's declarations imply are now raised whenever the bus is, and the grant is the pull.
2026-09-28 04:29:32 +02:00
mesh-admin f03e7b33c9 Merge pull request 'The controller may ask any module's tool' (#117) from fix/the-controller-may-ask-a-tool into main 2026-09-28 02:21:26 +00:00
6 changed files with 107 additions and 12 deletions
+25 -2
View File
@@ -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
View File
@@ -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
+21
View File
@@ -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
View File
@@ -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" }
} } } }
] ]
+14
View File
@@ -180,6 +180,20 @@ func (r Registry) MirrorImage(ctx context.Context, from, repository string) (str
if err != nil { if err != nil {
return "", err return "", err
} }
// **Already held is already mirrored.** A base is named by digest, and a digest this registry
// holds under the module's repository is the same bytes whatever upstream would say — so
// upstream is not asked. Asked every build, the public hub's anonymous pull limit was reached
// on the first merge that rebuilt a whole catalogue (2026-09-28), and every module whose base
// lives there failed on a copy it did not need.
if strings.HasPrefix(where.reference, "sha256:") {
held, err := r.has(ctx, "http://"+r.Address+"/v2/"+repository+"/manifests/"+where.reference)
if err != nil {
return "", fmt.Errorf("asking %s whether it holds %s: %w", r.Address, from, err)
}
if held {
return r.Address + "/" + repository + "@" + where.reference, nil
}
}
src := &source{client: r.client()} src := &source{client: r.client()}
digest, err := r.copyManifest(ctx, src, where, where.reference, repository) digest, err := r.copyManifest(ctx, src, where, where.reference, repository)
if err != nil { if err != nil {
+30
View File
@@ -103,6 +103,12 @@ func (m *theMeshsRegistry) handler() http.Handler {
m.mu.Lock() m.mu.Lock()
defer m.mu.Unlock() defer m.mu.Unlock()
switch { switch {
case r.Method == http.MethodHead && strings.Contains(r.URL.Path, "/manifests/"):
if _, ok := m.manifests[r.URL.Path[strings.LastIndex(r.URL.Path, "/")+1:]]; ok {
w.WriteHeader(http.StatusOK)
} else {
w.WriteHeader(http.StatusNotFound)
}
case r.Method == http.MethodHead && strings.Contains(r.URL.Path, "/blobs/"): case r.Method == http.MethodHead && strings.Contains(r.URL.Path, "/blobs/"):
if _, ok := m.blobs[r.URL.Path[strings.LastIndex(r.URL.Path, "/")+1:]]; ok { if _, ok := m.blobs[r.URL.Path[strings.LastIndex(r.URL.Path, "/")+1:]]; ok {
w.WriteHeader(http.StatusOK) w.WriteHeader(http.StatusOK)
@@ -219,3 +225,27 @@ func TestATagBeforeTheDigestIsNotPartOfTheRepository(t *testing.T) {
t.Fatalf("got %+v", got) t.Fatalf("got %+v", got)
} }
} }
// A base this registry already holds by digest is not asked of upstream at all: the public hub
// limits anonymous pulls, and a catalogue rebuilt on one merge asked it once per module.
func TestABaseAlreadyHeldIsNotAskedOfUpstream(t *testing.T) {
src, indexDigest, _ := anUpstreamRegistry(t)
dst := &theMeshsRegistry{blobs: map[string][]byte{}, manifests: map[string][]byte{}}
dstServer := httptest.NewServer(dst.handler())
defer dstServer.Close()
address := strings.TrimPrefix(dstServer.URL, "http://")
r := Registry{Address: address, HTTP: src.Client()}
host := strings.TrimPrefix(src.URL, "http://")
if _, err := r.MirrorImage(context.Background(), host+"/library/thing:latest", "hello-web/server"); err != nil {
t.Fatal(err)
}
// Upstream gone: the pinned base is answered from what the mesh holds.
src.Close()
reference, err := r.MirrorImage(context.Background(), host+"/library/thing@"+indexDigest, "hello-web/server")
if err != nil {
t.Fatalf("a base the registry holds was asked of an upstream that is gone: %v", err)
}
if reference != address+"/hello-web/server@"+indexDigest {
t.Fatalf("pinned as %q", reference)
}
}