Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
c5dc7e732a | ||
|
|
7efcccd013 | ||
|
|
964285f08c | ||
|
|
5698dda11f |
@@ -135,6 +135,16 @@ func run() error {
|
||||
// machine told about both would take work from one and answer on the other, and every log line would
|
||||
// say it was fine.
|
||||
func takeWorkFrom(credential Credential, on string) (link.BuildMachine, error) {
|
||||
// **The credential decides, before any variable does.** A machine moved to the new bus was
|
||||
// handed a credential for it and nothing else changed in its environment; that credential
|
||||
// names the bus by scheme, so it is enough to know which bus to take work from.
|
||||
if credential.onTheNewBus() {
|
||||
js, err := broker.DialPinned(credential.natsURL(), credential.Fingerprint)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return link.MachineOverNATS(js, on), nil
|
||||
}
|
||||
address, onNATS, err := broker.OnNATS()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
@@ -461,9 +471,25 @@ func brokerFrom() (Credential, error) {
|
||||
// out of band — here, sealed with the credential — and the endpoint is verified once at connect.
|
||||
type Credential struct {
|
||||
URL string `json:"url"`
|
||||
// Fingerprint is SHA-256 over the broker certificate's DER bytes, or empty to verify the
|
||||
// ordinary way.
|
||||
Fingerprint string `json:"fingerprint,omitempty"`
|
||||
// User and Password ride beside the address on the bus being built (design 25): a credential
|
||||
// embedded in a URL leaks into every log line that prints a connection, so the mesh seals them
|
||||
// as two fields and this machine joins them once, here, to dial.
|
||||
User string `json:"user,omitempty"`
|
||||
Password string `json:"password,omitempty"`
|
||||
}
|
||||
|
||||
// onTheNewBus is whether a credential is for the bus being built: its address says so, and the
|
||||
// mesh only ever seals such a credential with the user and password beside it.
|
||||
func (c Credential) onTheNewBus() bool { return strings.HasPrefix(strings.TrimSpace(c.URL), "nats://") }
|
||||
|
||||
// natsURL is the address with this machine's credential in it, for the one dial that needs it.
|
||||
func (c Credential) natsURL() string {
|
||||
rest := strings.TrimPrefix(strings.TrimSpace(c.URL), "nats://")
|
||||
if c.User == "" {
|
||||
return "nats://" + rest
|
||||
}
|
||||
return "nats://" + c.User + ":" + c.Password + "@" + rest
|
||||
}
|
||||
|
||||
// dial opens the connection, pinning the broker's certificate when there is one to pin.
|
||||
|
||||
@@ -761,10 +761,52 @@ func raiseTheBus(ctx context.Context, inv *inventory.Inventory, address string)
|
||||
// The work queues of the mesh's own roles (novox/hq ADR 0121). The queue before the holder,
|
||||
// deliberately: work queues until somebody arrives to do it, so assigning a build machine a week
|
||||
// after something started asking for builds flushes the backlog instead of having lost it.
|
||||
if err := broker.RaiseSeats(js, inventory.MeshSeats(), nil); err != nil {
|
||||
// With the seats' holders, so each role's work queue gets the consumer its holder takes
|
||||
// work from. Passed as nil until the first live raise, which left the build machine bound to a
|
||||
// consumer nothing had created (2026-09-28).
|
||||
holders, err := seatHolders(ctx, inv)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if err := broker.RaiseSeats(js, inventory.MeshSeats(), holders); err != nil {
|
||||
return err
|
||||
}
|
||||
fmt.Printf("the bus at %s has its streams, and %d machine(s) can hear a declaration\n",
|
||||
broker.BareAddress(address), len(names))
|
||||
return nil
|
||||
}
|
||||
|
||||
// seatHolders is who holds each of the mesh's seats, by seat name: the record where a handover
|
||||
// wrote one, and the assigned module claiming the seat otherwise — the same derivation the
|
||||
// resolver makes, read from the catalogue rather than re-resolved.
|
||||
func seatHolders(ctx context.Context, inv *inventory.Inventory) (map[string]broker.Holder, error) {
|
||||
out := map[string]broker.Holder{}
|
||||
entries, err := inv.Catalogued(ctx)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
for _, e := range entries {
|
||||
if len(e.On) == 0 {
|
||||
continue
|
||||
}
|
||||
for _, c := range e.Manifest.Claims {
|
||||
seat, known := catalogue.SeatNamed(c.Name)
|
||||
if !known {
|
||||
continue
|
||||
}
|
||||
if _, taken := out[seat.Name]; !taken {
|
||||
out[seat.Name] = broker.Holder{Node: e.On[0], Module: e.Manifest.Module}
|
||||
}
|
||||
}
|
||||
}
|
||||
recorded, err := inv.Holdings(ctx)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
for _, h := range recorded {
|
||||
if seat, known := catalogue.SeatNamed(h.Claim); known {
|
||||
out[seat.Name] = broker.Holder{Node: h.Node, Module: h.Module}
|
||||
}
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
|
||||
@@ -71,6 +71,20 @@ func pinnedTo(path string) (*tls.Config, error) {
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return PinnedToFingerprint(want), nil
|
||||
}
|
||||
|
||||
// DialPinned is Dial with the server's certificate pinned by a fingerprint the caller already holds
|
||||
// — a module or a build machine that was handed one beside its credential, and has no file.
|
||||
func DialPinned(url, fingerprint string, opts ...nats.Option) (*JetStream, error) {
|
||||
if strings.TrimSpace(fingerprint) != "" {
|
||||
opts = append(opts, nats.Secure(PinnedToFingerprint(fingerprint)))
|
||||
}
|
||||
return Dial(url, opts...)
|
||||
}
|
||||
|
||||
// PinnedToFingerprint accepts exactly the certificate with this SHA-256 and no other.
|
||||
func PinnedToFingerprint(want string) *tls.Config {
|
||||
return &tls.Config{
|
||||
InsecureSkipVerify: true, //nolint:gosec // replaced by the pin below, which is stricter
|
||||
MinVersion: tls.VersionTLS12,
|
||||
@@ -85,7 +99,7 @@ func pinnedTo(path string) (*tls.Config, error) {
|
||||
}
|
||||
return nil
|
||||
},
|
||||
}, nil
|
||||
}
|
||||
}
|
||||
|
||||
// Conn is the connection itself, for what the mesh keeps off JetStream on purpose — a heartbeat,
|
||||
|
||||
@@ -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))
|
||||
}
|
||||
|
||||
+1
-1
@@ -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" }
|
||||
} }
|
||||
|
||||
@@ -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",
|
||||
"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}",
|
||||
}
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user