Compare commits
3
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
28853a251b | ||
|
|
76a8b8df9e | ||
|
|
a3e8a4185b |
@@ -3,6 +3,7 @@ package main
|
|||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
"fmt"
|
"fmt"
|
||||||
|
"github.com/novox/mesh-controller/internal/broker"
|
||||||
"sort"
|
"sort"
|
||||||
"strings"
|
"strings"
|
||||||
|
|
||||||
@@ -67,6 +68,13 @@ func assign(ctx context.Context, open *stores, node, module string) (string, err
|
|||||||
for _, line := range settled {
|
for _, line := range settled {
|
||||||
said += "\n " + line
|
said += "\n " + line
|
||||||
}
|
}
|
||||||
|
// Its bus credential, in the same act (novox/hq issue 203): an assignment pushed before its
|
||||||
|
// credential exists delivers a process that cannot authenticate and crash-loops until somebody
|
||||||
|
// runs a second verb and a second push. Issued here when the module speaks on the bus and has
|
||||||
|
// no credential yet; kept when it has one, so re-assigning rotates nothing.
|
||||||
|
if line := issueOnAssign(ctx, open, node, module); line != "" {
|
||||||
|
said += "\n " + line
|
||||||
|
}
|
||||||
plan, _, err := planFor(ctx, open, node)
|
plan, _, err := planFor(ctx, open, node)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
// Kept, and still refused. Both halves are the answer, and the rest of the mesh is still
|
// Kept, and still refused. Both halves are the answer, and the rest of the mesh is still
|
||||||
@@ -152,3 +160,33 @@ func blockedElsewhere(ctx context.Context, open *stores, except string) string {
|
|||||||
out.WriteString("\nThis may or may not be what just changed — it is what is true now.")
|
out.WriteString("\nThis may or may not be what just changed — it is what is true now.")
|
||||||
return out.String()
|
return out.String()
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// issueOnAssign gives a newly assigned module its bus credential, the way `module issue` does, and
|
||||||
|
// says what it did in one line. Nothing for a module that declares no broker secret; nothing for one
|
||||||
|
// whose user is already minted (a credential is rotated on purpose, never by re-assigning); and when
|
||||||
|
// the bus cannot be reached from here, the line names the verb and the push that would refuse the
|
||||||
|
// module until it is run — never a silent placeholder (novox/hq issue 203).
|
||||||
|
func issueOnAssign(ctx context.Context, open *stores, node, module string) string {
|
||||||
|
inv := open.inventory
|
||||||
|
shelf, err := inv.Catalogue(ctx)
|
||||||
|
if err != nil {
|
||||||
|
return ""
|
||||||
|
}
|
||||||
|
m, known := shelf[module]
|
||||||
|
if !known || mayIssue(m) != nil {
|
||||||
|
return ""
|
||||||
|
}
|
||||||
|
user := broker.Principal{Kind: broker.KindModule, Node: node, Module: module}.Username()
|
||||||
|
if _, minted, err := inv.BusUserHash(ctx, user); err != nil || minted {
|
||||||
|
return ""
|
||||||
|
}
|
||||||
|
busAddress, err := broker.BusAddress()
|
||||||
|
if err == nil {
|
||||||
|
err = issueOnTheNewBus(ctx, inv, m, node, busAddress)
|
||||||
|
}
|
||||||
|
if err != nil {
|
||||||
|
return fmt.Sprintf("its bus credential is not issued (%v): `module issue %s --node %s` first — "+
|
||||||
|
"`push %s` refuses to send %s until it is", err, module, node, node, module)
|
||||||
|
}
|
||||||
|
return fmt.Sprintf("its bus credential is issued and sealed to %s, and arrives with the push", node)
|
||||||
|
}
|
||||||
|
|||||||
@@ -175,6 +175,12 @@ func TestTheControlPlaneIsToldWhereTheNodePutTheStoreAndTheBroker(t *testing.T)
|
|||||||
Guards: []int{15672},
|
Guards: []int{15672},
|
||||||
Resources: []map[string]any{{"id": "server", "type": "container", "name": "mesh-broker",
|
Resources: []map[string]any{{"id": "server", "type": "container", "name": "mesh-broker",
|
||||||
"ports": []any{"5671:5671", "5672:5672", "127.0.0.1:15672:15672"}, "image": "mq@" + aDigest}}})
|
"ports": []any{"5671:5671", "5672:5672", "127.0.0.1:15672:15672"}, "image": "mq@" + aDigest}}})
|
||||||
|
// The control plane's own bus user is the installer's, seeded at genesis before the controller
|
||||||
|
// runs (SeedBusUser); without it a push now refuses the credential nobody issued (issue 203).
|
||||||
|
if err := open.inventory.SeedBusUser(ctx, inventory.BusUser{Username: "anchor.mesh-controller",
|
||||||
|
Kind: inventory.BusController, Node: "anchor", Module: "mesh-controller"}, "bootstrap"); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
if _, err := assign(ctx, open, "anchor", "mesh-controller"); err != nil {
|
if _, err := assign(ctx, open, "anchor", "mesh-controller"); err != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -0,0 +1,90 @@
|
|||||||
|
package main
|
||||||
|
|
||||||
|
import (
|
||||||
|
"strings"
|
||||||
|
"testing"
|
||||||
|
|
||||||
|
"github.com/novox/mesh-controller/internal/catalogue"
|
||||||
|
"github.com/novox/mesh-controller/internal/inventory"
|
||||||
|
)
|
||||||
|
|
||||||
|
// A fresh assignment is pushed before its credential exists (novox/hq issue 203): `assign` recorded
|
||||||
|
// the module, `push` sealed a random own secret where the bus credential belongs, and the process
|
||||||
|
// crash-looped until a person ran `module issue` and pushed again. Now assigning a module that speaks
|
||||||
|
// on the bus issues its credential in the same act — or, when the bus cannot be reached from here,
|
||||||
|
// says which verb to run — and a push never seals a placeholder in a credential's place.
|
||||||
|
|
||||||
|
func aTalker() catalogue.Manifest {
|
||||||
|
return catalogue.Manifest{Module: "talker", Version: "1",
|
||||||
|
OwnSecrets: catalogue.OwnSecrets{"broker": {Path: "/var/lib/mesh/talker/broker"}},
|
||||||
|
Resources: []map[string]any{
|
||||||
|
{"id": "state", "type": "directory", "path": "/var/lib/mesh/talker", "mode": "0700"},
|
||||||
|
}}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestAssigningAModuleThatSpeaksOnTheBusNamesItsCredential(t *testing.T) {
|
||||||
|
open := aMesh(t)
|
||||||
|
ctx := t.Context()
|
||||||
|
register(t, open, aTalker())
|
||||||
|
|
||||||
|
// No bus is known to this process, so the credential cannot be issued here: the assignment
|
||||||
|
// stands and says exactly what must happen before a push — never silently.
|
||||||
|
said, err := assign(ctx, open, "laptop", "talker")
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if !strings.Contains(said, "module issue talker --node laptop") {
|
||||||
|
t.Fatalf("an assignment whose credential could not be issued does not name the verb:\n%s", said)
|
||||||
|
}
|
||||||
|
|
||||||
|
// And the push refuses to send it, naming the same verb, rather than sealing a placeholder.
|
||||||
|
plan, settings, err := planFor(ctx, open, "laptop")
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
_, err = declarationFor(ctx, open, "laptop", plan, settings)
|
||||||
|
if err == nil {
|
||||||
|
t.Fatal("a push sealed a placeholder where talker's bus credential belongs")
|
||||||
|
}
|
||||||
|
if !strings.Contains(err.Error(), "module issue talker --node laptop") || !strings.Contains(err.Error(), "issue 203") {
|
||||||
|
t.Fatalf("the refusal does not say what to run: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
// Once the user is minted, the push goes on to the credential the mesh sealed, and re-assigning
|
||||||
|
// does not mint again: a credential rotates on purpose, never by habit.
|
||||||
|
if _, err := open.inventory.MintBusPassword(ctx, inventory.BusUser{
|
||||||
|
Username: "laptop.talker", Kind: inventory.BusModule, Node: "laptop", Module: "talker"}); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
hash, _, err := open.inventory.BusUserHash(ctx, "laptop.talker")
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
said, err = assign(ctx, open, "laptop", "talker")
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if strings.Contains(said, "module issue") {
|
||||||
|
t.Fatalf("a module with a minted credential was told to issue one:\n%s", said)
|
||||||
|
}
|
||||||
|
again, _, err := open.inventory.BusUserHash(ctx, "laptop.talker")
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if again != hash {
|
||||||
|
t.Fatal("re-assigning rotated the credential")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// A module that declares no broker secret is left alone: nothing to issue, nothing said.
|
||||||
|
func TestAssigningAModuleThatDoesNotSpeakSaysNothingOfCredentials(t *testing.T) {
|
||||||
|
open := aMesh(t)
|
||||||
|
register(t, open, helloWeb())
|
||||||
|
said, err := assign(t.Context(), open, "laptop", "hello-web")
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if strings.Contains(said, "credential") {
|
||||||
|
t.Fatalf("a module without a broker secret was told about credentials:\n%s", said)
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -1,93 +0,0 @@
|
|||||||
package main
|
|
||||||
|
|
||||||
import (
|
|
||||||
"context"
|
|
||||||
"testing"
|
|
||||||
|
|
||||||
"github.com/novox/mesh-controller/internal/inventory"
|
|
||||||
)
|
|
||||||
|
|
||||||
// A declaration composed earlier is numbered lower than one composed later, whatever order the two
|
|
||||||
// are sent in (novox/hq issue 204). The number used to be taken at send time, after composing, so a
|
|
||||||
// declaration composed before an assignment changed and sent after a newer one carried the higher
|
|
||||||
// number — and the machine, which refuses a lower number, took the older content as the mesh's
|
|
||||||
// newest word. Taken before the composition reads anything, the order of numbers is the order of
|
|
||||||
// compositions, and the host's refusal does what it is for.
|
|
||||||
func TestADeclarationComposedEarlierIsNumberedLowerWhateverOrderItIsSent(t *testing.T) {
|
|
||||||
allot := numbered()
|
|
||||||
var composed []string
|
|
||||||
compose := func(stamp string) func(string) (sendable, error) {
|
|
||||||
return func(node string) (sendable, error) {
|
|
||||||
composed = append(composed, stamp)
|
|
||||||
return sendable{Resources: []map[string]any{{"id": node + "." + stamp}}}, nil
|
|
||||||
}
|
|
||||||
}
|
|
||||||
// Composed first — before an assignment changed — and sent last.
|
|
||||||
stale, _ := composeEach([]string{"anchor"}, allot, compose("before"))
|
|
||||||
// Composed after the change, sent first.
|
|
||||||
fresh, _ := composeEach([]string{"anchor"}, allot, compose("after"))
|
|
||||||
|
|
||||||
if stale[0].declared.Sequence != 1 || fresh[0].declared.Sequence != 2 {
|
|
||||||
t.Fatalf("the numbers do not follow the compositions: before=%d after=%d",
|
|
||||||
stale[0].declared.Sequence, fresh[0].declared.Sequence)
|
|
||||||
}
|
|
||||||
// Sent in the other order, the numbers do not change — so the machine that has applied the
|
|
||||||
// fresh one (2) refuses the stale one (1) when it arrives late.
|
|
||||||
if !(stale[0].declared.Sequence < fresh[0].declared.Sequence) {
|
|
||||||
t.Fatal("a declaration composed earlier must carry the lower number, however late it is sent")
|
|
||||||
}
|
|
||||||
if len(composed) != 2 || composed[0] != "before" {
|
|
||||||
t.Fatalf("compositions happened in an unexpected order: %v", composed)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// The number is taken before the first read of the composition, not after it: an allotter that
|
|
||||||
// fails leaves nothing composed for that machine, and the others are still composed.
|
|
||||||
func TestTheNumberIsTakenBeforeComposingAndItsFailureIsARefusal(t *testing.T) {
|
|
||||||
calls := 0
|
|
||||||
allot := func(node string) (int64, error) {
|
|
||||||
if node == "anchor" {
|
|
||||||
return 0, context.DeadlineExceeded
|
|
||||||
}
|
|
||||||
return 7, nil
|
|
||||||
}
|
|
||||||
sending, refusals := composeEach([]string{"anchor", "laptop"}, allot, func(node string) (sendable, error) {
|
|
||||||
calls++
|
|
||||||
if node == "anchor" {
|
|
||||||
t.Fatal("anchor was composed although its number could not be taken")
|
|
||||||
}
|
|
||||||
return sendable{}, nil
|
|
||||||
})
|
|
||||||
if calls != 1 || len(sending) != 1 || sending[0].node != "laptop" || sending[0].declared.Sequence != 7 {
|
|
||||||
t.Fatalf("laptop should be composed with its number and anchor refused: %v / %v", sending, refusals)
|
|
||||||
}
|
|
||||||
if len(refusals) != 1 {
|
|
||||||
t.Fatalf("anchor's failed number should be a refusal naming it: %v", refusals)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// What was sent is written down even when the sender's context is already cancelled (issue 204): a
|
|
||||||
// controller replaced mid-send had told the machine and never recorded it, so status read "applied,
|
|
||||||
// current" over a machine that had just been sent something else.
|
|
||||||
func TestASendIsRecordedEvenWhenTheSenderIsBeingCancelled(t *testing.T) {
|
|
||||||
inv := inventory.ForTest(t)
|
|
||||||
ctx, cancel := context.WithCancel(t.Context())
|
|
||||||
if _, err := inv.AddNode(ctx, "anchor"); err != nil {
|
|
||||||
t.Fatal(err)
|
|
||||||
}
|
|
||||||
cancel() // the sender is going away: its context is cancelled between the send and the record
|
|
||||||
body := []byte(`{"declaration":1,"resources":[]}`)
|
|
||||||
digest, err := recordSent(ctx, inv, "anchor", body)
|
|
||||||
if err != nil {
|
|
||||||
// NodeByName on the cancelled context may itself refuse; the record must still be possible
|
|
||||||
// through the detached context, so look the node up again on a live one.
|
|
||||||
t.Fatalf("recording a send after cancellation failed: %v", err)
|
|
||||||
}
|
|
||||||
outstanding, err := inv.Outstanding(t.Context(), "anchor")
|
|
||||||
if err != nil {
|
|
||||||
t.Fatal(err)
|
|
||||||
}
|
|
||||||
if outstanding != digest || digest != digestOf(body) {
|
|
||||||
t.Fatalf("the send was not recorded: outstanding %q, sent %q", outstanding, digest)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
@@ -533,6 +533,23 @@ func renderingFor(ctx context.Context, open *stores, node string,
|
|||||||
var sealed string
|
var sealed string
|
||||||
var err error
|
var err error
|
||||||
if choosing == Allocating {
|
if choosing == Allocating {
|
||||||
|
// **The broker credential is never invented here** (novox/hq issue 203). Every other
|
||||||
|
// own secret is the mesh's to make — a password nobody else knows — but this one
|
||||||
|
// is an account on the bus, minted by `module issue` and sealed by it; a push that
|
||||||
|
// made a random one would deliver a file the process cannot read and report the
|
||||||
|
// machine applied. Refused by name, with the verb.
|
||||||
|
if name == "broker" {
|
||||||
|
user := broker.Principal{Kind: broker.KindModule, Node: node, Module: m.Module}.Username()
|
||||||
|
if _, minted, err := inv.BusUserHash(ctx, user); err != nil {
|
||||||
|
return catalogue.Rendering{}, inventory.Node{}, err
|
||||||
|
} else if !minted {
|
||||||
|
return catalogue.Rendering{}, inventory.Node{}, fmt.Errorf(
|
||||||
|
"%s on %s has no bus credential: nothing was issued for %s, and a push "+
|
||||||
|
"would seal a placeholder its process cannot read (novox/hq issue 203). "+
|
||||||
|
"`module issue %s --node %s`, then push again",
|
||||||
|
m.Module, node, user, m.Module, node)
|
||||||
|
}
|
||||||
|
}
|
||||||
sealed, err = inv.SecretForModule(ctx, node, m.Module, name)
|
sealed, err = inv.SecretForModule(ctx, node, m.Module, name)
|
||||||
} else {
|
} else {
|
||||||
var held bool
|
var held bool
|
||||||
|
|||||||
+37
-65
@@ -205,12 +205,6 @@ func declare(ctx context.Context, args []string) error {
|
|||||||
if err := link.Declare(ctx, server.Bus(), ident, node, raw, 15*time.Second); err != nil {
|
if err := link.Declare(ctx, server.Bus(), ident, node, raw, 15*time.Second); err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
// Written down like every other send (novox/hq issue 204): a declaration a person sent by hand
|
|
||||||
// is still what the machine was last told, and status must not read it as current for the one
|
|
||||||
// the mesh would compose.
|
|
||||||
if _, err := recordSent(ctx, inv, node, raw); err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
fmt.Printf("sent %s a signed declaration (%d bytes)\n", node, len(raw))
|
fmt.Printf("sent %s a signed declaration (%d bytes)\n", node, len(raw))
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
@@ -354,7 +348,7 @@ func pushCommand(ctx context.Context, args []string) error {
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
sending, refusals := composeEach(asked, allotting(held, inv), func(node string) (sendable, error) {
|
sending, refusals := composeEach(asked, func(node string) (sendable, error) {
|
||||||
plan, settings, err := planFor(held, open, node)
|
plan, settings, err := planFor(held, open, node)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return sendable{}, err
|
return sendable{}, err
|
||||||
@@ -376,8 +370,12 @@ func pushCommand(ctx context.Context, args []string) error {
|
|||||||
sentDigest := map[string]string{}
|
sentDigest := map[string]string{}
|
||||||
defer release()
|
defer release()
|
||||||
for _, s := range sending {
|
for _, s := range sending {
|
||||||
// The number is inside the signed bytes, so a replayed older declaration cannot borrow a
|
// Numbered under the hold, one higher than the last, before the body exists — the number is
|
||||||
// newer one's (novox/hq 04-ISSUES/107); it was taken when the composition began (issue 204).
|
// inside the signed bytes, so a replayed older declaration cannot borrow a newer one's
|
||||||
|
// (novox/hq 04-ISSUES/107).
|
||||||
|
if err := number(ctx, inv, &s); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
body, err := s.declared.Body()
|
body, err := s.declared.Body()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
@@ -387,10 +385,14 @@ func pushCommand(ctx context.Context, args []string) error {
|
|||||||
}
|
}
|
||||||
// After it is away, not before. A digest recorded for something that failed to send would
|
// After it is away, not before. A digest recorded for something that failed to send would
|
||||||
// make the machine look current for a declaration it never received.
|
// make the machine look current for a declaration it never received.
|
||||||
digest, err := recordSent(ctx, inv, s.node, body)
|
record, err := inv.NodeByName(ctx, s.node)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
digest := digestOf(body)
|
||||||
|
if err := inv.RecordSent(ctx, record.ID, digest); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
sentDigest[s.node] = digest
|
sentDigest[s.node] = digest
|
||||||
fmt.Printf("sent %s %d resource(s)\n", s.node, len(s.declared.Resources))
|
fmt.Printf("sent %s %d resource(s)\n", s.node, len(s.declared.Resources))
|
||||||
}
|
}
|
||||||
@@ -474,7 +476,11 @@ func pushCommand(ctx context.Context, args []string) error {
|
|||||||
15*time.Second); err != nil {
|
15*time.Second); err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
if _, err := recordSent(ctx, inv, s.node, body); err != nil {
|
record, err := inv.NodeByName(ctx, s.node)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
if err := inv.RecordSent(ctx, record.ID, digestOf(body)); err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
fmt.Printf("sent %s %d resource(s)\n", s.node, len(s.declared.Resources))
|
fmt.Printf("sent %s %d resource(s)\n", s.node, len(s.declared.Resources))
|
||||||
@@ -563,31 +569,17 @@ type readyNode struct {
|
|||||||
//
|
//
|
||||||
// The all-or-nothing rule is kept where it means something — sendTo, which rotates a credential
|
// The all-or-nothing rule is kept where it means something — sendTo, which rotates a credential
|
||||||
// across two machines that must agree — and dropped here, where it never did.
|
// across two machines that must agree — and dropped here, where it never did.
|
||||||
func composeEach(names []string, allot func(node string) (int64, error),
|
func composeEach(names []string,
|
||||||
compose func(node string) (sendable, error)) ([]readyNode, []string) {
|
compose func(node string) (sendable, error)) ([]readyNode, []string) {
|
||||||
|
|
||||||
var sending []readyNode
|
var sending []readyNode
|
||||||
var refusals []string
|
var refusals []string
|
||||||
for _, name := range names {
|
for _, name := range names {
|
||||||
// **Numbered before it is composed, not before it is sent** (novox/hq issue 204). The
|
|
||||||
// number says where this declaration stands against every other the mesh composed for the
|
|
||||||
// machine, and the host refuses one lower than the last it applied. Taken at send time, as
|
|
||||||
// it was, a declaration composed a minute ago — before an assignment changed — went out with
|
|
||||||
// a number higher than one composed after the change and sent before it, and the machine
|
|
||||||
// took the older content as the newer word: on 2026-10-02 a runtime assigned and applied on
|
|
||||||
// two machines was undone two seconds later by exactly that. Taken here, before the first
|
|
||||||
// read, what was composed earlier is numbered lower whatever order the sends happen in.
|
|
||||||
seq, err := allot(name)
|
|
||||||
if err != nil {
|
|
||||||
refusals = append(refusals, fmt.Sprintf("%s:\n%v", name, err))
|
|
||||||
continue
|
|
||||||
}
|
|
||||||
declared, err := compose(name)
|
declared, err := compose(name)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
refusals = append(refusals, fmt.Sprintf("%s:\n%v", name, err))
|
refusals = append(refusals, fmt.Sprintf("%s:\n%v", name, err))
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
declared.Sequence = seq
|
|
||||||
if len(declared.Resources) == 0 {
|
if len(declared.Resources) == 0 {
|
||||||
// Sent, not skipped (novox/hq issue 127). A node whose declaration composes to
|
// Sent, not skipped (novox/hq issue 127). A node whose declaration composes to
|
||||||
// nothing may have HELD something before — the broker opening a placement gave it,
|
// nothing may have HELD something before — the broker opening a placement gave it,
|
||||||
@@ -615,10 +607,13 @@ func sendRound(ctx context.Context, open *stores, names []string,
|
|||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
defer release()
|
defer release()
|
||||||
sending, refused := composeEach(names, allotting(held, open.inventory), func(node string) (sendable, error) {
|
sending, refused := composeEach(names, func(node string) (sendable, error) {
|
||||||
return compose(held, node)
|
return compose(held, node)
|
||||||
})
|
})
|
||||||
for _, s := range sending {
|
for _, s := range sending {
|
||||||
|
if err := number(ctx, open.inventory, &s); err != nil {
|
||||||
|
return refused, err
|
||||||
|
}
|
||||||
body, err := s.declared.Body()
|
body, err := s.declared.Body()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return refused, err
|
return refused, err
|
||||||
@@ -675,12 +670,6 @@ func sendTo(ctx context.Context, open *stores, names []string) error {
|
|||||||
var sending []readyNode
|
var sending []readyNode
|
||||||
var refusals []string
|
var refusals []string
|
||||||
for _, name := range names {
|
for _, name := range names {
|
||||||
// Numbered before composing, for the reason composeEach gives (novox/hq issue 204).
|
|
||||||
seq, err := allot(ctx, inv, name)
|
|
||||||
if err != nil {
|
|
||||||
refusals = append(refusals, fmt.Sprintf("%s:\n%v", name, err))
|
|
||||||
continue
|
|
||||||
}
|
|
||||||
plan, settings, err := planFor(ctx, open, name)
|
plan, settings, err := planFor(ctx, open, name)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
refusals = append(refusals, fmt.Sprintf("%s:\n%v", name, err))
|
refusals = append(refusals, fmt.Sprintf("%s:\n%v", name, err))
|
||||||
@@ -692,7 +681,6 @@ func sendTo(ctx context.Context, open *stores, names []string) error {
|
|||||||
refusals = append(refusals, fmt.Sprintf("%s:\n%v", name, err))
|
refusals = append(refusals, fmt.Sprintf("%s:\n%v", name, err))
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
declared.Sequence = seq
|
|
||||||
reportLeftOut(name, declared)
|
reportLeftOut(name, declared)
|
||||||
sending = append(sending, readyNode{name, declared})
|
sending = append(sending, readyNode{name, declared})
|
||||||
}
|
}
|
||||||
@@ -708,6 +696,9 @@ func sendTo(ctx context.Context, open *stores, names []string) error {
|
|||||||
defer server.Close()
|
defer server.Close()
|
||||||
|
|
||||||
for _, s := range sending {
|
for _, s := range sending {
|
||||||
|
if err := number(ctx, inv, &s); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
body, err := s.declared.Body()
|
body, err := s.declared.Body()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
@@ -715,7 +706,11 @@ func sendTo(ctx context.Context, open *stores, names []string) error {
|
|||||||
if err := link.Declare(ctx, server.Bus(), ident, s.node, body, 15*time.Second); err != nil {
|
if err := link.Declare(ctx, server.Bus(), ident, s.node, body, 15*time.Second); err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
if _, err := recordSent(ctx, inv, s.node, body); err != nil {
|
record, err := inv.NodeByName(ctx, s.node)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
if err := inv.RecordSent(ctx, record.ID, digestOf(body)); err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
fmt.Printf(" sent %s %d resource(s)\n", s.node, len(s.declared.Resources))
|
fmt.Printf(" sent %s %d resource(s)\n", s.node, len(s.declared.Resources))
|
||||||
@@ -956,38 +951,15 @@ func seatHolders(ctx context.Context, inv *inventory.Inventory) (map[string]brok
|
|||||||
}
|
}
|
||||||
|
|
||||||
// number gives one send the next sequence for its node (novox/hq 04-ISSUES/107).
|
// number gives one send the next sequence for its node (novox/hq 04-ISSUES/107).
|
||||||
// allotting is allot over one inventory, in the shape composeEach takes.
|
func number(ctx context.Context, inv *inventory.Inventory, s *readyNode) error {
|
||||||
func allotting(ctx context.Context, inv *inventory.Inventory) func(node string) (int64, error) {
|
record, err := inv.NodeByName(ctx, s.node)
|
||||||
return func(node string) (int64, error) { return allot(ctx, inv, node) }
|
|
||||||
}
|
|
||||||
|
|
||||||
// allot takes the next sequence for a machine — the number its next declaration carries.
|
|
||||||
func allot(ctx context.Context, inv *inventory.Inventory, node string) (int64, error) {
|
|
||||||
record, err := inv.NodeByName(ctx, node)
|
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return 0, err
|
return err
|
||||||
}
|
}
|
||||||
return inv.NextSequence(ctx, record.ID)
|
seq, err := inv.NextSequence(ctx, record.ID)
|
||||||
}
|
|
||||||
|
|
||||||
// recordSent writes down what a machine was just sent, and returns the digest.
|
|
||||||
//
|
|
||||||
// **On a context that outlives the caller's** (novox/hq issue 204). The record is written after the
|
|
||||||
// declaration is away, so a send that failed is never recorded as current — and a controller being
|
|
||||||
// replaced mid-send had its context cancelled between the two, so the machine was told and the mesh
|
|
||||||
// never wrote it down: status read "applied, current" over a machine that had just been sent
|
|
||||||
// something else. What was sent was sent; the record of it must not depend on the sender living
|
|
||||||
// another second. Bounded, so a store that is away does not hold a dying process open for ever.
|
|
||||||
func recordSent(ctx context.Context, inv *inventory.Inventory, node string, body []byte) (string, error) {
|
|
||||||
kept, cancel := context.WithTimeout(context.WithoutCancel(ctx), 10*time.Second)
|
|
||||||
defer cancel()
|
|
||||||
record, err := inv.NodeByName(kept, node)
|
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return "", err
|
return err
|
||||||
}
|
}
|
||||||
digest := digestOf(body)
|
s.declared.Sequence = seq
|
||||||
if err := inv.RecordSent(kept, record.ID, digest); err != nil {
|
return nil
|
||||||
return "", err
|
|
||||||
}
|
|
||||||
return digest, nil
|
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -17,7 +17,7 @@ import (
|
|||||||
// the wrong machine no longer refuses the whole node), applied one level up.
|
// the wrong machine no longer refuses the whole node), applied one level up.
|
||||||
func TestOneUnresolvableNodeStillLetsTheRestBeSent(t *testing.T) {
|
func TestOneUnresolvableNodeStillLetsTheRestBeSent(t *testing.T) {
|
||||||
sending, refusals := composeEach(
|
sending, refusals := composeEach(
|
||||||
[]string{"anchor", "home-server", "laptop"}, numbered(),
|
[]string{"anchor", "home-server", "laptop"},
|
||||||
func(node string) (sendable, error) {
|
func(node string) (sendable, error) {
|
||||||
if node == "anchor" {
|
if node == "anchor" {
|
||||||
return sendable{}, errors.New(`nothing provides "acme-ca", wanted by route-proxy`)
|
return sendable{}, errors.New(`nothing provides "acme-ca", wanted by route-proxy`)
|
||||||
@@ -43,7 +43,7 @@ func TestOneUnresolvableNodeStillLetsTheRestBeSent(t *testing.T) {
|
|||||||
// (novox/hq issue 127): it may have held something before, and only sending the empty
|
// (novox/hq issue 127): it may have held something before, and only sending the empty
|
||||||
// declaration tells it to drop what the mesh owned. It is never a refusal.
|
// declaration tells it to drop what the mesh owned. It is never a refusal.
|
||||||
func TestAnEmptyDeclarationIsSentSoTheNodeDropsWhatItHeld(t *testing.T) {
|
func TestAnEmptyDeclarationIsSentSoTheNodeDropsWhatItHeld(t *testing.T) {
|
||||||
sending, refusals := composeEach([]string{"spare"}, numbered(),
|
sending, refusals := composeEach([]string{"spare"},
|
||||||
func(string) (sendable, error) { return sendable{}, nil })
|
func(string) (sendable, error) { return sendable{}, nil })
|
||||||
if len(sending) != 1 || len(refusals) != 0 {
|
if len(sending) != 1 || len(refusals) != 0 {
|
||||||
t.Errorf("an empty declaration must be sent, not skipped or refused: %v / %v", sending, refusals)
|
t.Errorf("an empty declaration must be sent, not skipped or refused: %v / %v", sending, refusals)
|
||||||
@@ -74,9 +74,3 @@ func TestASkippedMachineIsStillAnError(t *testing.T) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// numbered is an allotter for tests: one higher per call, as the inventory's is per machine.
|
|
||||||
func numbered() func(string) (int64, error) {
|
|
||||||
var n int64
|
|
||||||
return func(string) (int64, error) { n++; return n, nil }
|
|
||||||
}
|
|
||||||
|
|||||||
@@ -36,17 +36,29 @@ func tiersOf(set []string, edges []inventory.Edge) [][]string {
|
|||||||
for _, m := range set {
|
for _, m := range set {
|
||||||
deps[m] = map[string]bool{}
|
deps[m] = map[string]bool{}
|
||||||
}
|
}
|
||||||
|
// The build seat's holders follow the controller that defines their worker (EdgeWorkerOf,
|
||||||
|
// novox/hq issue 206), so the built-by edge from that controller to such a holder yields: the
|
||||||
|
// controller is built by whichever build machine is running, as the runtime image always was.
|
||||||
|
worker := map[string]map[string]bool{}
|
||||||
|
for _, e := range edges {
|
||||||
|
if e.Kind == inventory.EdgeWorkerOf && in[e.From] && in[e.To] {
|
||||||
|
if worker[e.To] == nil {
|
||||||
|
worker[e.To] = map[string]bool{}
|
||||||
|
}
|
||||||
|
worker[e.To][e.From] = true
|
||||||
|
}
|
||||||
|
}
|
||||||
for _, e := range edges {
|
for _, e := range edges {
|
||||||
// A code dependency — B packages A's source — rebuilds B with A, in the same tier: B's
|
// A code dependency — B packages A's source — rebuilds B with A, in the same tier: B's
|
||||||
// build needs nothing of A's first. The other kinds order: stands-on and declared after
|
// build needs nothing of A's first. The other kinds order: stands-on and declared after
|
||||||
// the base is built, built-by after the build machine is built and running — except for
|
// the base is built, built-by after the build machine is built and running — except for
|
||||||
// what the build machine itself stands on. The runtime image is built by the builder and
|
// what the build machine itself stands on, and for the controller whose worker the build
|
||||||
// the builder is built on the runtime image; the image comes first, built by the builder
|
// machine binds. The runtime image is built by the builder and the builder is built on the
|
||||||
// that is running, which is the only one there could be.
|
// runtime image; the image comes first, built by the builder that is running.
|
||||||
if !in[e.From] || !in[e.To] || e.From == e.To || e.Kind == inventory.EdgePackages {
|
if !in[e.From] || !in[e.To] || e.From == e.To || e.Kind == inventory.EdgePackages {
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
if e.Kind == inventory.EdgeBuiltBy && isBaseOf(e.From, e.To, edges, in) {
|
if e.Kind == inventory.EdgeBuiltBy && (isBaseOf(e.From, e.To, edges, in) || worker[e.From][e.To]) {
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
deps[e.From][e.To] = true
|
deps[e.From][e.To] = true
|
||||||
@@ -122,7 +134,10 @@ func reachableFrom(moved []string, edges []inventory.Edge) []string {
|
|||||||
for grew := true; grew; {
|
for grew := true; grew; {
|
||||||
grew = false
|
grew = false
|
||||||
for _, e := range edges {
|
for _, e := range edges {
|
||||||
if e.Kind == inventory.EdgeBuiltBy {
|
// Built-by and worker-of order a plan; neither widens it. A new build machine changes
|
||||||
|
// nothing it builds, and a new controller changes nothing about the holder it orders —
|
||||||
|
// what packages the controller's source is already a code edge.
|
||||||
|
if e.Kind == inventory.EdgeBuiltBy || e.Kind == inventory.EdgeWorkerOf {
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
if in[e.To] && !in[e.From] {
|
if in[e.To] && !in[e.From] {
|
||||||
|
|||||||
@@ -0,0 +1,45 @@
|
|||||||
|
package main
|
||||||
|
|
||||||
|
import (
|
||||||
|
"testing"
|
||||||
|
|
||||||
|
"github.com/novox/mesh-controller/internal/inventory"
|
||||||
|
)
|
||||||
|
|
||||||
|
// The holder of the build seat follows the controller that defines its worker (novox/hq issue 206).
|
||||||
|
// On 2026-10-03 a plan put the build machine in tier 0 and the controller in tier 1; the new build
|
||||||
|
// machine could not bind the worker the old controller had defined, and nothing could build the
|
||||||
|
// controller that would have redefined it. The built-by edge from the controller to its build
|
||||||
|
// machine yields to that order: the controller is built by whichever build machine is running.
|
||||||
|
func TestTheBuildSeatsHolderFollowsTheControllerThatDefinesItsWorker(t *testing.T) {
|
||||||
|
edges := []inventory.Edge{
|
||||||
|
{From: "build-agent", To: "mesh-controller", Kind: inventory.EdgePackages},
|
||||||
|
{From: "build-agent", To: "mesh-controller", Kind: inventory.EdgeWorkerOf},
|
||||||
|
{From: "mesh-controller", To: "build-agent", Kind: inventory.EdgeBuiltBy},
|
||||||
|
{From: "route-proxy", To: "mesh-controller", Kind: inventory.EdgePackages},
|
||||||
|
{From: "route-proxy", To: "build-agent", Kind: inventory.EdgeBuiltBy},
|
||||||
|
}
|
||||||
|
set := reachableFrom([]string{"mesh-controller"}, edges)
|
||||||
|
if len(set) != 3 {
|
||||||
|
t.Fatalf("the controller, what packages it, and nothing more: %v", set)
|
||||||
|
}
|
||||||
|
tiers := tiersOf(set, edges)
|
||||||
|
pos := map[string]int{}
|
||||||
|
for i, tier := range tiers {
|
||||||
|
for _, m := range tier {
|
||||||
|
pos[m] = i
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if pos["mesh-controller"] != 0 {
|
||||||
|
t.Fatalf("the controller first, built by the build machine that is running: %v", tiers)
|
||||||
|
}
|
||||||
|
if pos["build-agent"] <= pos["mesh-controller"] {
|
||||||
|
t.Fatalf("the build machine after the controller that defines its worker: %v", tiers)
|
||||||
|
}
|
||||||
|
if pos["route-proxy"] <= pos["build-agent"] {
|
||||||
|
t.Fatalf("what the build machine builds comes after it: %v", tiers)
|
||||||
|
}
|
||||||
|
if hasCycle(tiers, edges) {
|
||||||
|
t.Fatalf("no cycle here: %v", tiers)
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -214,6 +214,43 @@ func (j *JetStream) EnsureConsumer(c Consumer) error {
|
|||||||
|
|
||||||
switch have, err := j.js.ConsumerInfo(c.Stream, c.Name); {
|
switch have, err := j.js.ConsumerInfo(c.Stream, c.Name); {
|
||||||
case err == nil:
|
case err == nil:
|
||||||
|
// **The controller owns the worker's shape, type included** (novox/hq issue 206). A holder
|
||||||
|
// built for a pull worker cannot bind a push one — `cannot pull subscribe to push based
|
||||||
|
// consumer` — and on 2026-10-03 the build machine rolled before the controller that would
|
||||||
|
// have redefined its worker, restarted on that for an hour, and nothing could build the
|
||||||
|
// controller that would have ended it. The server cannot change a consumer's type in place,
|
||||||
|
// so one of the wrong type is re-made: on a work queue nothing is lost, because what was
|
||||||
|
// acknowledged is gone from the stream and what was not is delivered again from the start.
|
||||||
|
// On any other stream a re-made consumer would replay what this one acknowledged (issue
|
||||||
|
// 156), so there it is said and left, and the person re-makes it knowing the cost.
|
||||||
|
if havePush, wantPush := have.Config.DeliverSubject != "", want.DeliverSubject != ""; havePush != wantPush {
|
||||||
|
shape := func(push bool) string {
|
||||||
|
if push {
|
||||||
|
return "push"
|
||||||
|
}
|
||||||
|
return "pull"
|
||||||
|
}
|
||||||
|
info, err := j.js.StreamInfo(c.Stream)
|
||||||
|
if err != nil {
|
||||||
|
return fmt.Errorf("asking about stream %s to re-make consumer %s: %w", c.Stream, c.Name, err)
|
||||||
|
}
|
||||||
|
if info.Config.Retention != nats.WorkQueuePolicy {
|
||||||
|
j.note("consumer %s on %s is %s and should be %s; not re-made, because %s keeps its history "+
|
||||||
|
"and a re-made consumer replays what this one acknowledged (novox/hq issue 156). Re-make it by hand",
|
||||||
|
c.Name, c.Stream, shape(havePush), shape(wantPush), c.Stream)
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
j.note("consumer %s on %s changes from %s to %s delivery: re-made where it left off, nothing "+
|
||||||
|
"acknowledged comes back and nothing pending is lost (novox/hq issue 206); a holder bound to "+
|
||||||
|
"the old shape binds again", c.Name, c.Stream, shape(havePush), shape(wantPush))
|
||||||
|
if err := j.js.DeleteConsumer(c.Stream, c.Name); err != nil {
|
||||||
|
return fmt.Errorf("re-making consumer %s on %s as %s: %w", c.Name, c.Stream, shape(wantPush), err)
|
||||||
|
}
|
||||||
|
if _, err := j.js.AddConsumer(c.Stream, want); err != nil {
|
||||||
|
return fmt.Errorf("re-making consumer %s on %s as %s: %w", c.Name, c.Stream, shape(wantPush), err)
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
// Where an existing consumer starts is its history, not something an assertion may move:
|
// Where an existing consumer starts is its history, not something an assertion may move:
|
||||||
// the server refuses a changed deliver policy outright. Carried across, so asserting twice
|
// the server refuses a changed deliver policy outright. Carried across, so asserting twice
|
||||||
// is the no-op a restart depends on.
|
// is the no-op a restart depends on.
|
||||||
|
|||||||
@@ -0,0 +1,145 @@
|
|||||||
|
package broker
|
||||||
|
|
||||||
|
import (
|
||||||
|
"os"
|
||||||
|
"testing"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"github.com/nats-io/nats.go"
|
||||||
|
)
|
||||||
|
|
||||||
|
// A seat's worker that changed from push to pull delivery strands a holder built for the new shape
|
||||||
|
// (novox/hq issue 206): the server refuses a pull subscription on a push consumer, and the controller
|
||||||
|
// that would redefine it was the build that nobody could take. The controller owns the worker's
|
||||||
|
// shape, type included: on a work queue it re-makes one of the wrong type, losing nothing, and a
|
||||||
|
// pull subscription then binds and takes what was pending.
|
||||||
|
//
|
||||||
|
// docker run -d --rm --name t -p 14231:4222 nats:2.10-alpine -js
|
||||||
|
// MESH_TEST_NATS=nats://127.0.0.1:14231 go test ./internal/broker/ -run TestAWorker
|
||||||
|
func TestAWorkerOfTheWrongTypeIsRemadeOnAWorkQueueAndAPullThenBinds(t *testing.T) {
|
||||||
|
url := os.Getenv("MESH_TEST_NATS")
|
||||||
|
if url == "" {
|
||||||
|
t.Skip("MESH_TEST_NATS unset")
|
||||||
|
}
|
||||||
|
js, err := Dial(url)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
defer js.Close()
|
||||||
|
|
||||||
|
const stream, worker, filter = "SEAT_T_SHELF", "SEAT_T_SHELF_worker", "mesh.seat.t-shelf.accept.>"
|
||||||
|
_ = js.js.DeleteStream(stream)
|
||||||
|
if _, err := js.js.AddStream(&nats.StreamConfig{
|
||||||
|
Name: stream, Subjects: []string{filter}, Retention: nats.WorkQueuePolicy, Storage: nats.MemoryStorage,
|
||||||
|
}); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
defer func() { _ = js.js.DeleteStream(stream) }()
|
||||||
|
|
||||||
|
// The worker as the previous controller defined it: push, in a queue group.
|
||||||
|
if _, err := js.js.AddConsumer(stream, &nats.ConsumerConfig{
|
||||||
|
Durable: worker, AckPolicy: nats.AckExplicitPolicy, AckWait: 60 * time.Second, MaxDeliver: 5,
|
||||||
|
FilterSubject: filter, DeliverSubject: "_DELIVER." + worker, DeliverGroup: "holders",
|
||||||
|
}); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
for _, body := range []string{"one", "two", "three"} {
|
||||||
|
if _, err := js.js.Publish("mesh.seat.t-shelf.accept.build", []byte(body)); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
// The old holder took and acknowledged the first ask, then went away.
|
||||||
|
old, err := js.js.QueueSubscribeSync(filter, "holders", nats.Bind(stream, worker))
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
m, err := old.NextMsg(twoSeconds)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if string(m.Data) != "one" {
|
||||||
|
t.Fatalf("the first ask is %q", m.Data)
|
||||||
|
}
|
||||||
|
if err := m.AckSync(); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if err := old.Unsubscribe(); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
|
||||||
|
// The new controller asserts the worker as the mesh derives it now: pull.
|
||||||
|
if err := js.EnsureConsumer(Consumer{
|
||||||
|
Name: worker, Stream: stream, Filters: []string{filter}, AckWaitSeconds: 60, MaxDeliver: 5,
|
||||||
|
Why: "the test's worker",
|
||||||
|
}); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
have, err := js.js.ConsumerInfo(stream, worker)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if have.Config.DeliverSubject != "" || have.Config.DeliverGroup != "" {
|
||||||
|
t.Fatalf("the worker is still push: %+v", have.Config)
|
||||||
|
}
|
||||||
|
|
||||||
|
// A holder built for the new shape binds, and takes exactly what the old one left.
|
||||||
|
sub, err := js.js.PullSubscribe(filter, worker, nats.Bind(stream, worker), nats.ManualAck())
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("a pull subscription does not bind the re-made worker: %v", err)
|
||||||
|
}
|
||||||
|
got, err := sub.Fetch(3, nats.MaxWait(twoSeconds))
|
||||||
|
if err != nil && len(got) == 0 {
|
||||||
|
t.Fatalf("nothing pending was delivered: %v", err)
|
||||||
|
}
|
||||||
|
var bodies []string
|
||||||
|
for _, g := range got {
|
||||||
|
bodies = append(bodies, string(g.Data))
|
||||||
|
_ = g.Ack()
|
||||||
|
}
|
||||||
|
if len(bodies) != 2 || bodies[0] != "two" || bodies[1] != "three" {
|
||||||
|
t.Fatalf("the pending asks after the acknowledged one, in order: %v", bodies)
|
||||||
|
}
|
||||||
|
|
||||||
|
// Asserted again, the pull worker is the no-op a restart depends on.
|
||||||
|
if err := js.EnsureConsumer(Consumer{
|
||||||
|
Name: worker, Stream: stream, Filters: []string{filter}, AckWaitSeconds: 60, MaxDeliver: 5,
|
||||||
|
}); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// On a stream that keeps its history, a worker of the wrong type is said and left: re-making it would
|
||||||
|
// replay what it acknowledged (novox/hq issue 156), and that is a person's call.
|
||||||
|
func TestAWorkerOfTheWrongTypeOnAHistoryStreamIsLeftAndSaid(t *testing.T) {
|
||||||
|
url := os.Getenv("MESH_TEST_NATS")
|
||||||
|
if url == "" {
|
||||||
|
t.Skip("MESH_TEST_NATS unset")
|
||||||
|
}
|
||||||
|
js, err := Dial(url)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
defer js.Close()
|
||||||
|
const stream, worker, filter = "EVENTS_T", "EVENTS_T_reader", "mesh.t.event.>"
|
||||||
|
_ = js.js.DeleteStream(stream)
|
||||||
|
if _, err := js.js.AddStream(&nats.StreamConfig{Name: stream, Subjects: []string{filter}, Storage: nats.MemoryStorage}); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
defer func() { _ = js.js.DeleteStream(stream) }()
|
||||||
|
if _, err := js.js.AddConsumer(stream, &nats.ConsumerConfig{
|
||||||
|
Durable: worker, AckPolicy: nats.AckExplicitPolicy, AckWait: 60 * time.Second,
|
||||||
|
FilterSubject: filter, DeliverSubject: "_DELIVER." + worker,
|
||||||
|
}); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if err := js.EnsureConsumer(Consumer{Name: worker, Stream: stream, Filters: []string{filter}, AckWaitSeconds: 60}); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
have, err := js.js.ConsumerInfo(stream, worker)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if have.Config.DeliverSubject == "" {
|
||||||
|
t.Fatal("a history stream's consumer was re-made, which replays what it acknowledged")
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -872,6 +872,15 @@ func (r Resolution) compose(with Rendering, owner map[string]string,
|
|||||||
if renamed := reflectsRenamed(m.Module, resource["reload-on"]); renamed != nil {
|
if renamed := reflectsRenamed(m.Module, resource["reload-on"]); renamed != nil {
|
||||||
copied["reload-on"] = renamed
|
copied["reload-on"] = renamed
|
||||||
}
|
}
|
||||||
|
// **What reads one of this module's own secrets is restarted when it changes** (novox/hq
|
||||||
|
// issue 203, issue 206). A credential is re-issued by the mesh, and a container that
|
||||||
|
// mounted the old file keeps the old one open: the build machine ran for an hour on a
|
||||||
|
// credential the mesh had replaced, because its manifest restarted it on its
|
||||||
|
// environment file and nobody had thought to name the credential too. Composed here so
|
||||||
|
// no manifest has to say it, for a container or a daemon that names the secret's path.
|
||||||
|
if reads := secretsReadBy(copied, m); len(reads) > 0 {
|
||||||
|
copied["restart-on"] = withRestartOn(copied["restart-on"], reads)
|
||||||
|
}
|
||||||
// **A version prepares its state before it runs** (novox/hq ADR 0135). Derived from the
|
// **A version prepares its state before it runs** (novox/hq ADR 0135). Derived from the
|
||||||
// module's own resource rather than declared beside it: what prepares the state is the
|
// module's own resource rather than declared beside it: what prepares the state is the
|
||||||
// module's own code, so what it is given has to be what that code is given — and a
|
// module's own code, so what it is given has to be what that code is given — and a
|
||||||
@@ -2079,3 +2088,87 @@ func portOfEndpoint(values map[string]any, ports map[string]int) {
|
|||||||
values["port"] = port
|
values["port"] = port
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// secretsReadBy is the file resources of this module's own secrets that a container or a daemon reads
|
||||||
|
// — named in its volumes, its environment or its env-files by the secret's placed path — as
|
||||||
|
// restart-on ids. Nothing for other shapes, and nothing for a scheduled or run-once process, which
|
||||||
|
// the host refuses a restart-on for (it runs again anyway, and reads the file afresh).
|
||||||
|
func secretsReadBy(resource map[string]any, m Manifest) []string {
|
||||||
|
kind := fmt.Sprint(resource["type"])
|
||||||
|
if kind != "container" && kind != "process" {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
if resource["schedule"] != nil || resource["run-once"] == true {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
var mentioned []string
|
||||||
|
for _, key := range []string{"volumes", "env", "env-file"} {
|
||||||
|
mentioned = append(mentioned, stringsIn(resource[key])...)
|
||||||
|
}
|
||||||
|
var out []string
|
||||||
|
for _, name := range sortedKeys(m.OwnSecrets) {
|
||||||
|
path := m.OwnSecrets[name].Path
|
||||||
|
if path == "" {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
for _, s := range mentioned {
|
||||||
|
// A volume is `source:destination[:mode]`; an env value or an env-file is the path itself.
|
||||||
|
if s == path || strings.HasPrefix(s, path+":") {
|
||||||
|
out = append(out, m.Module+"."+NeedID(name))
|
||||||
|
break
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return out
|
||||||
|
}
|
||||||
|
|
||||||
|
// stringsIn is every string in a list or a map's values; nothing for anything else.
|
||||||
|
func stringsIn(v any) []string {
|
||||||
|
switch x := v.(type) {
|
||||||
|
case []any:
|
||||||
|
var out []string
|
||||||
|
for _, item := range x {
|
||||||
|
if s, ok := item.(string); ok {
|
||||||
|
out = append(out, s)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return out
|
||||||
|
case []string:
|
||||||
|
return x
|
||||||
|
case map[string]any:
|
||||||
|
var out []string
|
||||||
|
for _, k := range sortedKeys(x) {
|
||||||
|
if s, ok := x[k].(string); ok {
|
||||||
|
out = append(out, s)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return out
|
||||||
|
case map[string]string:
|
||||||
|
var out []string
|
||||||
|
for _, k := range sortedKeys(x) {
|
||||||
|
out = append(out, x[k])
|
||||||
|
}
|
||||||
|
return out
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// withRestartOn is a resource's restart-on list with these ids added once each.
|
||||||
|
func withRestartOn(have any, add []string) []any {
|
||||||
|
var out []any
|
||||||
|
seen := map[string]bool{}
|
||||||
|
for _, id := range reflectsRenamed("", have) {
|
||||||
|
s := fmt.Sprint(id)
|
||||||
|
if !seen[s] {
|
||||||
|
seen[s] = true
|
||||||
|
out = append(out, s)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
for _, id := range add {
|
||||||
|
if !seen[id] {
|
||||||
|
seen[id] = true
|
||||||
|
out = append(out, id)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return out
|
||||||
|
}
|
||||||
|
|||||||
@@ -0,0 +1,52 @@
|
|||||||
|
package catalogue
|
||||||
|
|
||||||
|
import (
|
||||||
|
"reflect"
|
||||||
|
"strings"
|
||||||
|
"testing"
|
||||||
|
)
|
||||||
|
|
||||||
|
// What reads one of a module's own secrets is restarted when the secret changes (novox/hq issue 203,
|
||||||
|
// issue 206): the build machine kept an hour-old credential open because its manifest restarted it
|
||||||
|
// on its environment file alone. Composed, so a manifest need not say it; a scheduled process is
|
||||||
|
// left alone, because the host refuses a restart-on for one and it reads the file afresh each run.
|
||||||
|
func TestAContainerReadingAnOwnSecretIsRestartedWhenItChanges(t *testing.T) {
|
||||||
|
m := Manifest{Module: "agent", Version: "1",
|
||||||
|
OwnSecrets: OwnSecrets{"broker": {Path: "/var/lib/mesh/agent/broker"}},
|
||||||
|
Resources: []map[string]any{
|
||||||
|
{"id": "mesh-state", "type": "directory", "path": "/var/lib/mesh/agent", "mode": "0700"},
|
||||||
|
{"id": "settings", "type": "file", "path": "/var/lib/mesh/agent/agent.env", "mode": "0600", "content": "A=1\n"},
|
||||||
|
{"id": "server", "type": "container", "name": "agent", "network": "host",
|
||||||
|
"image": "registry.example/agent@sha256:" + strings.Repeat("a", 64),
|
||||||
|
"volumes": []any{"/var/lib/mesh/agent:/run/mesh:ro", "/var/lib/mesh/agent/broker:/run/mesh/broker:ro"},
|
||||||
|
"env-file": []any{"/var/lib/mesh/agent/agent.env"},
|
||||||
|
"restart-on": []any{"settings"}},
|
||||||
|
{"id": "nightly", "type": "container", "name": "agent-nightly", "schedule": "0 3 * * *",
|
||||||
|
"image": "registry.example/agent@sha256:" + strings.Repeat("a", 64),
|
||||||
|
"volumes": []any{"/var/lib/mesh/agent/broker:/run/mesh/broker:ro"}},
|
||||||
|
{"id": "other", "type": "container", "name": "agent-other",
|
||||||
|
"image": "registry.example/agent@sha256:" + strings.Repeat("a", 64)},
|
||||||
|
}}
|
||||||
|
got, err := Resolve(shelf(m), []string{m.Module},
|
||||||
|
Node{Name: "anchor", At: "10.0.0.1", Capabilities: map[string]bool{"container-runtime": true}}, World{})
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
out, err := got.Declaration(Rendering{Needed: map[string]map[string]string{"agent": {"broker": "SEALED"}}})
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
by := map[string]map[string]any{}
|
||||||
|
for _, r := range out {
|
||||||
|
by[r["id"].(string)] = r
|
||||||
|
}
|
||||||
|
if want := []any{"agent.settings", "agent.needs-broker"}; !reflect.DeepEqual(by["agent.server"]["restart-on"], want) {
|
||||||
|
t.Fatalf("the server reads the credential and is not restarted on it: %v", by["agent.server"]["restart-on"])
|
||||||
|
}
|
||||||
|
if _, has := by["agent.nightly"]["restart-on"]; has {
|
||||||
|
t.Fatalf("a scheduled container was given a restart-on, which the host refuses: %v", by["agent.nightly"]["restart-on"])
|
||||||
|
}
|
||||||
|
if _, has := by["agent.other"]["restart-on"]; has {
|
||||||
|
t.Fatalf("a container that reads no secret was given one to restart on: %v", by["agent.other"]["restart-on"])
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -18,8 +18,17 @@ const (
|
|||||||
EdgeBuiltBy = "built-by"
|
EdgeBuiltBy = "built-by"
|
||||||
// EdgeDeclared: the manifest's own `build.on`.
|
// EdgeDeclared: the manifest's own `build.on`.
|
||||||
EdgeDeclared = "declared"
|
EdgeDeclared = "declared"
|
||||||
|
// EdgeWorkerOf: the module holds the build seat, whose worker the control plane defines
|
||||||
|
// (novox/hq issue 206). The one place a *running* order enters the graph: a build machine rolled
|
||||||
|
// before the controller that redefines its worker cannot bind it, and nothing can then build the
|
||||||
|
// controller that would end that — so the holder of the build seat follows the controller, and
|
||||||
|
// the controller is built by whichever build machine is running, as it always was.
|
||||||
|
EdgeWorkerOf = "worker-of"
|
||||||
)
|
)
|
||||||
|
|
||||||
|
// TheControlPlane is the module that defines every seat's worker on the bus.
|
||||||
|
const TheControlPlane = "mesh-controller"
|
||||||
|
|
||||||
// Edge is one dependency: From depends on To, in the way Kind says.
|
// Edge is one dependency: From depends on To, in the way Kind says.
|
||||||
type Edge struct {
|
type Edge struct {
|
||||||
From string `json:"from"`
|
From string `json:"from"`
|
||||||
@@ -106,6 +115,13 @@ func dependenciesOf(entries []Entry, against map[string][]string, read map[strin
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
if known[TheControlPlane] {
|
||||||
|
for _, b := range builders {
|
||||||
|
if b != TheControlPlane {
|
||||||
|
add(b, TheControlPlane, EdgeWorkerOf)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
sort.Slice(out, func(a, b int) bool {
|
sort.Slice(out, func(a, b int) bool {
|
||||||
if out[a].From != out[b].From {
|
if out[a].From != out[b].From {
|
||||||
return out[a].From < out[b].From
|
return out[a].From < out[b].From
|
||||||
|
|||||||
Reference in New Issue
Block a user