The build machine takes work on the bus its credential names, and the work queue has a taker
Two halves of one gap the first build over the new bus met. The machine decided its bus from a variable its container never received, so the credential the mesh sealed to it went unread; a credential for the new bus names the bus by scheme and carries user, password and fingerprint beside the address, and that is enough to dial it, pinned. And the roles' work queues were raised with no holders, so the consumer a machine binds to take work was never created: the holders are read from the catalogue and the handover record, as the resolver reads them.
This commit is contained in:
@@ -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.
|
||||||
|
|||||||
@@ -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)
|
||||||
|
|||||||
@@ -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
|
||||||
|
}
|
||||||
|
|||||||
@@ -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,
|
||||||
|
|||||||
Reference in New Issue
Block a user