Merge pull request 'The build machine takes work on the bus its credential names, and the work queue has a taker' (#107) from feat/the-build-machine-takes-work-on-nats into main

This commit was merged in pull request #107.
This commit is contained in:
2026-09-27 23:59:56 +00:00
4 changed files with 89 additions and 7 deletions
+29 -3
View File
@@ -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 // machine told about both would take work from one and answer on the other, and every log line would
// say it was fine. // say it was fine.
func takeWorkFrom(credential Credential, on string) (link.BuildMachine, error) { 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() address, onNATS, err := broker.OnNATS()
if err != nil { if err != nil {
return nil, err return nil, err
@@ -460,10 +470,26 @@ func brokerFrom() (Credential, error) {
// **The same shape a node gets, for the same reason** (novox/hq ADR 0004): the fingerprint travels // **The same shape a node gets, for the same reason** (novox/hq ADR 0004): the fingerprint travels
// out of band — here, sealed with the credential — and the endpoint is verified once at connect. // out of band — here, sealed with the credential — and the endpoint is verified once at connect.
type Credential struct { type Credential struct {
URL string `json:"url"` 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"` 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. // dial opens the connection, pinning the broker's certificate when there is one to pin.
+2 -2
View File
@@ -14,8 +14,8 @@ func TestBuilderDiagnosticsStayOffStdout(t *testing.T) {
allowed := map[string]bool{ allowed := map[string]bool{
"string(body)": true, // once.go: the result JSON, which IS stdout "string(body)": true, // once.go: the result JSON, which IS stdout
"version)": true, // --version "version)": true, // --version
`"stopping")`: true, // the loop.s shutdown line `"stopping")`: true, // the loop.s shutdown line
"usage)": true, // --help text, for a human "usage)": true, // --help text, for a human
} }
for _, file := range []string{"once.go", "main.go"} { for _, file := range []string{"once.go", "main.go"} {
src, err := os.ReadFile(file) src, err := os.ReadFile(file)
+43 -1
View File
@@ -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, // 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 // 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. // 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 return err
} }
fmt.Printf("the bus at %s has its streams, and %d machine(s) can hear a declaration\n", fmt.Printf("the bus at %s has its streams, and %d machine(s) can hear a declaration\n",
broker.BareAddress(address), len(names)) broker.BareAddress(address), len(names))
return nil 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
}
+15 -1
View File
@@ -71,6 +71,20 @@ func pinnedTo(path string) (*tls.Config, error) {
if err != nil { if err != nil {
return nil, err 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{ return &tls.Config{
InsecureSkipVerify: true, //nolint:gosec // replaced by the pin below, which is stricter InsecureSkipVerify: true, //nolint:gosec // replaced by the pin below, which is stricter
MinVersion: tls.VersionTLS12, MinVersion: tls.VersionTLS12,
@@ -85,7 +99,7 @@ func pinnedTo(path string) (*tls.Config, error) {
} }
return nil return nil
}, },
}, nil }
} }
// Conn is the connection itself, for what the mesh keeps off JetStream on purpose — a heartbeat, // Conn is the connection itself, for what the mesh keeps off JetStream on purpose — a heartbeat,