Compare commits
13
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
ce9e20fbbc | ||
|
|
878690697e | ||
|
|
ad97297576 | ||
|
|
683b1ed693 | ||
|
|
04f9f378b0 | ||
|
|
1c3f44a526 | ||
|
|
89e152dfe2 | ||
|
|
1ebad3786c | ||
|
|
f2f526a60a | ||
|
|
4b4c7e0e0d | ||
|
|
cec792ce9d | ||
|
|
338d033632 | ||
|
|
1be926cec4 |
@@ -88,7 +88,7 @@ func reportsReaching(t *testing.T, open *stores, reachable []link.Reach, held ..
|
||||
if err := open.inventory.RecordSent(ctx, record.ID, digestOf(body)); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := (link.Enrolment{Inventory: open.inventory}).Heard(ctx, link.Report{
|
||||
if _, err := (link.Enrolment{Inventory: open.inventory}).Heard(ctx, link.Report{
|
||||
Node: "anchor", Applied: []string{"hello-web.x"}, Declared: digestOf(body),
|
||||
Firewall: "ufw", Held: held, Reachable: reachable,
|
||||
}); err != nil {
|
||||
|
||||
@@ -162,7 +162,13 @@ func declare(ctx context.Context, args []string) error {
|
||||
return err
|
||||
}
|
||||
|
||||
server, err := connectLink(ctx, nil, nil, nil)
|
||||
// **With the inventory, so the bus is raised** (novox/hq ADR 0134, design 30). A module's
|
||||
// declaration and how it hears what it consumes move together: its consumer is derived from the
|
||||
// same records this declaration is composed from. Raised only when the control plane started
|
||||
// serving, a module that gained a `consumes` was sent a declaration it could act on and a
|
||||
// consumer that never delivered the event — and nothing anywhere said the two disagreed
|
||||
// (found on review, 2026-09-28). Everything the raise does is idempotent.
|
||||
server, err := connectLink(ctx, inv, nil, nil)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
@@ -181,6 +181,14 @@ func PermissionsFor(p Principal) (Permissions, error) {
|
||||
for _, seat := range meshSeatsTheControllerUses {
|
||||
pub = append(pub, "mesh.seat."+seat+".accept.>")
|
||||
}
|
||||
// **And what the mesh says it did** (novox/hq ADR 0134). The control plane states its own
|
||||
// facts under the seat it holds, because a role's events belong to the role and keep their
|
||||
// address while the holder is replaced. Named one by one rather than as a whole namespace:
|
||||
// least authority, and a fact nothing states is authority nobody uses.
|
||||
for _, event := range ControllerStates {
|
||||
pub = append(pub, seatEventSubject(ControllerSeat, event))
|
||||
}
|
||||
|
||||
// Every module's tools: **the control plane is the way in** (novox/hq ADR 0095). A person
|
||||
// or an agent asks through it and every question passes one process where an audit
|
||||
// belongs — so it, alone among principals, may call any tool by name. The first `ask` on
|
||||
|
||||
@@ -0,0 +1,44 @@
|
||||
package broker_test
|
||||
|
||||
import (
|
||||
"slices"
|
||||
"testing"
|
||||
|
||||
"github.com/novox/mesh-controller/internal/broker"
|
||||
"github.com/novox/mesh-controller/internal/catalogue"
|
||||
"github.com/novox/mesh-controller/internal/link"
|
||||
)
|
||||
|
||||
// The facts the control plane states are named twice — in the grant that permits them and in the code
|
||||
// that states them — because `link` imports `broker` and the dependency cannot go the other way. So a
|
||||
// test keeps them agreeing: a subject the grant omits is refused at the moment the mesh has something
|
||||
// to say, and one the grant adds that nothing states is authority nobody uses.
|
||||
//
|
||||
// An external test package, because it may import both while neither imports the other.
|
||||
func TestTheFactsTheGrantPermitsAreTheFactsTheMeshStates(t *testing.T) {
|
||||
if broker.ControllerSeat != link.MeshControllerSeat {
|
||||
t.Fatalf("the grant is written for the %q seat and the mesh states its facts under %q",
|
||||
broker.ControllerSeat, link.MeshControllerSeat)
|
||||
}
|
||||
for _, event := range []string{link.KeyApplied, link.KeyRefused, link.KeyBuiltBefore} {
|
||||
if !slices.Contains(broker.ControllerStates, event) {
|
||||
t.Errorf("the mesh states %q and its account may not publish it", event)
|
||||
}
|
||||
}
|
||||
if len(broker.ControllerStates) != 3 {
|
||||
t.Errorf("the grant permits %v, which is more than the mesh states", broker.ControllerStates)
|
||||
}
|
||||
// **And the seat says it.** A seat carries the protocol of its role (novox/hq ADR 0129), so the
|
||||
// facts the control plane states are the seat's `emits` — which is what lets anything else declare
|
||||
// that it consumes them, and what the subject-agreement check reads to know they have an owner.
|
||||
var declared []string
|
||||
for _, seat := range catalogue.SeatsWithAProtocol() {
|
||||
if seat.Name == broker.ControllerSeat {
|
||||
declared = seat.Emits
|
||||
}
|
||||
}
|
||||
if !slices.Equal(declared, broker.ControllerStates) {
|
||||
t.Errorf("the %s seat emits %v and the grant permits %v", broker.ControllerSeat,
|
||||
declared, broker.ControllerStates)
|
||||
}
|
||||
}
|
||||
@@ -160,6 +160,17 @@ func Overlaps() []string {
|
||||
// ack subject is derived from (nats.go: `$JS.ACK.<stream>.controller.>`).
|
||||
const ControllerName = "controller"
|
||||
|
||||
// ControllerSeat is the role the control plane holds, and ControllerStates are the facts it states
|
||||
// under it (novox/hq ADR 0134).
|
||||
//
|
||||
// **Written here as well as in `link`, and a test keeps them agreeing.** `link` imports `broker`, so
|
||||
// `broker` cannot import `link`; a grant naming a subject the controller never publishes is authority
|
||||
// nobody uses, and a controller publishing one the grant omits is refused at the moment it has
|
||||
// something to say.
|
||||
const ControllerSeat = "mesh-controller"
|
||||
|
||||
var ControllerStates = []string{"applied", "refused", "built-before"}
|
||||
|
||||
// ControllerFollows are the events the controller reacts to: the catalogue saying a module's
|
||||
// current version moved, and a catalogue that has just started saying it may have missed builds.
|
||||
//
|
||||
|
||||
+1
-1
@@ -24,7 +24,7 @@ accounts {
|
||||
jetstream: enabled
|
||||
users = [
|
||||
{ user: "controller", password: "$2a$11$cccccccccccccccccccccc", permissions: {
|
||||
publish: { allow: ["$JS.ACK.CONTROL.controller.>", "$JS.ACK.EVENTS.controller.>", "$JS.API.>", "_INBOX.enrol.>", "mesh.control.>", "mesh.mod.*.tool.>", "mesh.node.>", "mesh.seat.mesh-build-machine.accept.>"] }
|
||||
publish: { allow: ["$JS.ACK.CONTROL.controller.>", "$JS.ACK.EVENTS.controller.>", "$JS.API.>", "_INBOX.enrol.>", "mesh.control.>", "mesh.mod.*.tool.>", "mesh.node.>", "mesh.seat.mesh-build-machine.accept.>", "mesh.seat.mesh-controller.event.applied", "mesh.seat.mesh-controller.event.built-before", "mesh.seat.mesh-controller.event.refused"] }
|
||||
subscribe: { allow: ["$JS.API.>", "_DELIVER.controller", "_DELIVER.controller.>", "_INBOX.controller.>", "mesh.control.>", "mesh.mod.gitea.event.pull.merged", "mesh.mod.mesh-catalog.event.catching-up", "mesh.mod.mesh-catalog.event.upgraded", "mesh.seat.mesh-build-machine.event.built"] }
|
||||
allow_responses: { max: 1, ttl: "1m" }
|
||||
} }
|
||||
|
||||
@@ -186,7 +186,7 @@ func (r Registry) MirrorImage(ctx context.Context, from, repository string) (str
|
||||
// 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)
|
||||
held, err := r.has(ctx, "http://"+r.Address+"/v2/"+repository+"/manifests/"+where.reference, manifestAccept)
|
||||
if err != nil {
|
||||
return "", fmt.Errorf("asking %s whether it holds %s: %w", r.Address, from, err)
|
||||
}
|
||||
|
||||
@@ -104,6 +104,14 @@ func (m *theMeshsRegistry) handler() http.Handler {
|
||||
defer m.mu.Unlock()
|
||||
switch {
|
||||
case r.Method == http.MethodHead && strings.Contains(r.URL.Path, "/manifests/"):
|
||||
// **As strictly as a real registry.** A manifest is answered only in a media type the
|
||||
// caller named; a request with no Accept is answered as if nothing were there. The fake
|
||||
// used to answer regardless, which is why it could not catch a check that asked without
|
||||
// one — and the mesh copied every base again (2026-09-28).
|
||||
if !strings.Contains(r.Header.Get("Accept"), "manifest") && !strings.Contains(r.Header.Get("Accept"), "index") {
|
||||
w.WriteHeader(http.StatusNotFound)
|
||||
return
|
||||
}
|
||||
if _, ok := m.manifests[r.URL.Path[strings.LastIndex(r.URL.Path, "/")+1:]]; ok {
|
||||
w.WriteHeader(http.StatusOK)
|
||||
} else {
|
||||
|
||||
@@ -122,11 +122,23 @@ func (r Registry) PublishArchive(ctx context.Context, repository string, body []
|
||||
return final, nil
|
||||
}
|
||||
|
||||
func (r Registry) has(ctx context.Context, url string) (bool, error) {
|
||||
// has is whether this registry already holds what is at that URL.
|
||||
//
|
||||
// **A manifest HEAD must say what it accepts.** A registry answers a manifest request only in a media
|
||||
// type the caller named, and a bare HEAD — no Accept at all — is answered 404 for a manifest it holds
|
||||
// perfectly well. Measured against the mesh's own registry (2026-09-28): the same digest answered 200
|
||||
// with the manifest media types and 404 without them, so a check written without them concluded the
|
||||
// registry held nothing, copied every base again, and exhausted the public hub's pull limit. A blob
|
||||
// needs no Accept, which is why this went unnoticed: the same helper was right for blobs and wrong
|
||||
// for manifests.
|
||||
func (r Registry) has(ctx context.Context, url string, accept ...string) (bool, error) {
|
||||
request, err := http.NewRequestWithContext(ctx, http.MethodHead, url, nil)
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
for _, media := range accept {
|
||||
request.Header.Add("Accept", media)
|
||||
}
|
||||
response, err := r.client().Do(request)
|
||||
if err != nil {
|
||||
return false, fmt.Errorf("cannot reach the registry at %s: %w", r.Address, err)
|
||||
|
||||
@@ -1685,7 +1685,10 @@ func prepared(from map[string]any) map[string]any {
|
||||
for k, v := range from {
|
||||
step[k] = v
|
||||
}
|
||||
step["id"] = fmt.Sprint(from["id"]) + ".prepare"
|
||||
// **A hyphen, not a dot.** A resource's id is `<module>.<its own id>`, and a module's name may
|
||||
// itself contain a dot (`novox.be`), so the module is everything before the *last* dot — which
|
||||
// only works if what the mesh derives adds no dot of its own.
|
||||
step["id"] = fmt.Sprint(from["id"]) + "-prepare"
|
||||
step["name"] = fmt.Sprint(from["name"]) + "-prepare"
|
||||
step["run-once"] = true
|
||||
step["args"] = []any{PreparationArgument}
|
||||
|
||||
@@ -3,6 +3,7 @@ package catalogue
|
||||
import (
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"strings"
|
||||
"testing"
|
||||
)
|
||||
|
||||
@@ -57,13 +58,18 @@ func TestThePreparationRunsTheModulesOwnCodeAndComesRightBeforeIt(t *testing.T)
|
||||
ids := idsOf(out)
|
||||
at := -1
|
||||
for i, id := range ids {
|
||||
if id == "gitea.runtime.prepare" {
|
||||
if id == "gitea.runtime-prepare" {
|
||||
at = i
|
||||
}
|
||||
}
|
||||
if at < 0 {
|
||||
t.Fatalf("nothing prepares this module's state: %v", ids)
|
||||
}
|
||||
// A module's name may contain a dot, so a resource's module is everything before the last one —
|
||||
// which the derived id must not add to, or a machine reads the wrong owner from it.
|
||||
if strings.Count("gitea.runtime-prepare", ".") != 1 {
|
||||
t.Fatal("the derived id adds a dot, so what owns it cannot be read from it")
|
||||
}
|
||||
if ids[at+1] != "gitea.runtime" {
|
||||
t.Fatalf("the preparation is not immediately before the module's own code: %v", ids)
|
||||
}
|
||||
@@ -79,7 +85,7 @@ func TestThePreparationRunsTheModulesOwnCodeAndComesRightBeforeIt(t *testing.T)
|
||||
func TestThePreparationIsGivenWhatTheModuleIsGiven(t *testing.T) {
|
||||
out := declaredFor(t, aPreparingModule())
|
||||
declared := byID(out)
|
||||
step, workload := declared["gitea.runtime.prepare"], declared["gitea.runtime"]
|
||||
step, workload := declared["gitea.runtime-prepare"], declared["gitea.runtime"]
|
||||
if step == nil || workload == nil {
|
||||
t.Fatalf("expected both, got %v", idsOf(out))
|
||||
}
|
||||
@@ -108,7 +114,7 @@ func TestAModuleThatPreparesNothingGetsNoStep(t *testing.T) {
|
||||
m := aPreparingModule()
|
||||
m.Prepares = false
|
||||
for _, id := range idsOf(declaredFor(t, m)) {
|
||||
if id == "gitea.runtime.prepare" {
|
||||
if id == "gitea.runtime-prepare" {
|
||||
t.Fatal("a module that prepares nothing was given a preparation")
|
||||
}
|
||||
}
|
||||
|
||||
@@ -48,7 +48,11 @@ type Seat struct {
|
||||
//
|
||||
// In the order a person reads it: the mesh's own, then a node's.
|
||||
var defaultSeats = []Seat{
|
||||
{Name: "mesh-controller", Scope: ScopeMesh, Decision: "novox/hq ADR 0079"},
|
||||
// The control plane states what it did under the seat it holds (novox/hq ADR 0134): a role's
|
||||
// events belong to the role, so they keep their address while the holder is replaced. No accepts,
|
||||
// so no work queue is raised for it — only what its holder may say.
|
||||
{Name: "mesh-controller", Scope: ScopeMesh, Decision: "novox/hq ADR 0079",
|
||||
Emits: []string{"applied", "refused", "built-before"}},
|
||||
{Name: "mesh-store", Scope: ScopeMesh, Delivers: "postgres-database", Decision: "novox/hq ADR 0079"},
|
||||
// **Delivers the mesh's own bus, not `amqp`.** Those were the same word until
|
||||
// ADR 0127 separated them: `amqp` is a backing service a module may require, and this seat is
|
||||
|
||||
@@ -30,12 +30,12 @@ func TestRefusedAndFailedAreDifferentSituations(t *testing.T) {
|
||||
refuser := nodeNamed(t, inv, "refuser")
|
||||
failer := nodeNamed(t, inv, "failer")
|
||||
|
||||
if err := inv.RecordDoing(ctx, refuser, Doing{
|
||||
if _, err := inv.RecordDoing(ctx, refuser, Doing{
|
||||
Outcome: OutcomeRefused, Refused: "resource \"x\": a file needs a path",
|
||||
}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := inv.RecordDoing(ctx, failer, Doing{
|
||||
if _, err := inv.RecordDoing(ctx, failer, Doing{
|
||||
Outcome: OutcomeFailed,
|
||||
Failed: []FailedResource{{ID: "svc", Error: "unit not found"}},
|
||||
Applied: 4,
|
||||
@@ -70,7 +70,7 @@ func TestAMachineDoingWhatItWasToldIsNotOnTheList(t *testing.T) {
|
||||
inv := fresh(t)
|
||||
ctx := context.Background()
|
||||
id := nodeNamed(t, inv, "fine")
|
||||
if err := inv.RecordDoing(ctx, id, Doing{Outcome: OutcomeApplied, Applied: 6}); err != nil {
|
||||
if _, err := inv.RecordDoing(ctx, id, Doing{Outcome: OutcomeApplied, Applied: 6}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
wrong, err := inv.NotDoingWhatTheyWereTold(ctx)
|
||||
@@ -97,12 +97,12 @@ func TestTheLastReportReplacesTheOneBefore(t *testing.T) {
|
||||
inv := fresh(t)
|
||||
ctx := context.Background()
|
||||
id := nodeNamed(t, inv, "recovered")
|
||||
if err := inv.RecordDoing(ctx, id, Doing{
|
||||
if _, err := inv.RecordDoing(ctx, id, Doing{
|
||||
Outcome: OutcomeFailed, Failed: []FailedResource{{ID: "a", Error: "no"}},
|
||||
}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := inv.RecordDoing(ctx, id, Doing{Outcome: OutcomeApplied, Applied: 3}); err != nil {
|
||||
if _, err := inv.RecordDoing(ctx, id, Doing{Outcome: OutcomeApplied, Applied: 3}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
wrong, err := inv.NotDoingWhatTheyWereTold(ctx)
|
||||
@@ -141,7 +141,7 @@ func TestWhatANodeSaidGoesWhenTheNodeDoes(t *testing.T) {
|
||||
inv := fresh(t)
|
||||
ctx := context.Background()
|
||||
id := nodeNamed(t, inv, "leaving")
|
||||
if err := inv.RecordDoing(ctx, id, Doing{Outcome: OutcomeFailed}); err != nil {
|
||||
if _, err := inv.RecordDoing(ctx, id, Doing{Outcome: OutcomeFailed}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := inv.store.Pool().Exec(ctx, `delete from node where name = 'leaving'`); err != nil {
|
||||
@@ -257,7 +257,7 @@ func TestTheSameFailureReportedAgainIsCountedNotRestarted(t *testing.T) {
|
||||
id := nodeNamed(t, inv, "looping")
|
||||
same := Doing{Outcome: OutcomeFailed, Failed: []FailedResource{{ID: "img", Error: "no such image"}}}
|
||||
|
||||
if err := inv.RecordDoing(ctx, id, same); err != nil {
|
||||
if _, err := inv.RecordDoing(ctx, id, same); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
first, _, err := inv.DoingOf(ctx, "looping")
|
||||
@@ -269,7 +269,7 @@ func TestTheSameFailureReportedAgainIsCountedNotRestarted(t *testing.T) {
|
||||
}
|
||||
|
||||
for range StuckAfter - 1 {
|
||||
if err := inv.RecordDoing(ctx, id, same); err != nil {
|
||||
if _, err := inv.RecordDoing(ctx, id, same); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
@@ -287,7 +287,7 @@ func TestTheSameFailureReportedAgainIsCountedNotRestarted(t *testing.T) {
|
||||
// The same resource failing with different words — a duration, a counter — is still the same
|
||||
// failure: it is the resource that loops, not the sentence.
|
||||
reworded := Doing{Outcome: OutcomeFailed, Failed: []FailedResource{{ID: "img", Error: "no such image (after 31s)"}}}
|
||||
if err := inv.RecordDoing(ctx, id, reworded); err != nil {
|
||||
if _, err := inv.RecordDoing(ctx, id, reworded); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
still, _, err := inv.DoingOf(ctx, "looping")
|
||||
@@ -300,7 +300,7 @@ func TestTheSameFailureReportedAgainIsCountedNotRestarted(t *testing.T) {
|
||||
|
||||
// A different failure is a new situation, not a longer one.
|
||||
other := Doing{Outcome: OutcomeFailed, Failed: []FailedResource{{ID: "svc", Error: "unit not found"}}}
|
||||
if err := inv.RecordDoing(ctx, id, other); err != nil {
|
||||
if _, err := inv.RecordDoing(ctx, id, other); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
changed, _, err := inv.DoingOf(ctx, "looping")
|
||||
@@ -312,7 +312,7 @@ func TestTheSameFailureReportedAgainIsCountedNotRestarted(t *testing.T) {
|
||||
}
|
||||
|
||||
// And a clean apply clears it: the machine is doing what it was told, since nothing.
|
||||
if err := inv.RecordDoing(ctx, id, Doing{Outcome: OutcomeApplied, Applied: 2}); err != nil {
|
||||
if _, err := inv.RecordDoing(ctx, id, Doing{Outcome: OutcomeApplied, Applied: 2}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
fine, _, err := inv.DoingOf(ctx, "looping")
|
||||
@@ -324,7 +324,7 @@ func TestTheSameFailureReportedAgainIsCountedNotRestarted(t *testing.T) {
|
||||
}
|
||||
|
||||
// The list of what is wrong carries the count, so `status` can say it.
|
||||
if err := inv.RecordDoing(ctx, id, same); err != nil {
|
||||
if _, err := inv.RecordDoing(ctx, id, same); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
wrong, err := inv.NotDoingWhatTheyWereTold(ctx)
|
||||
|
||||
+29
-18
@@ -670,31 +670,39 @@ func sameFailure(a, b Doing) bool {
|
||||
// a clean apply clears both (novox/hq 04-ISSUES/065). The previous row is read first and the
|
||||
// comparison made here, so "the same" is a rule this package states rather than a jsonb equality
|
||||
// that would restart the count on a changed word in an error.
|
||||
func (i *Inventory) RecordDoing(ctx context.Context, node string, d Doing) error {
|
||||
// **And whether this report was news**, which is what makes a fact about it worth stating (novox/hq
|
||||
// ADR 0134). A machine reconciles continuously and reports each time; the same outcome about the same
|
||||
// declaration is the same state said again, and a fact per report would be a fact per minute per
|
||||
// machine that tells nobody anything. Read here because the previous row is read here anyway.
|
||||
func (i *Inventory) RecordDoing(ctx context.Context, node string, d Doing) (news bool, err error) {
|
||||
failed, err := json.Marshal(d.Failed)
|
||||
if err != nil {
|
||||
return err
|
||||
return false, err
|
||||
}
|
||||
var before Doing
|
||||
var beforeFailed []byte
|
||||
found := i.store.Pool().QueryRow(ctx,
|
||||
`select outcome, refused, failed, failing_since, failures, coalesce(declared,'')
|
||||
from node_report where node = $1`,
|
||||
node).Scan(&before.Outcome, &before.Refused, &beforeFailed, &before.Since, &before.Times,
|
||||
&before.Declared)
|
||||
switch {
|
||||
case errors.Is(found, pgx.ErrNoRows):
|
||||
news = true
|
||||
case found != nil:
|
||||
return false, found
|
||||
default:
|
||||
if err := json.Unmarshal(beforeFailed, &before.Failed); err != nil {
|
||||
return false, err
|
||||
}
|
||||
news = before.Outcome != d.Outcome || before.Declared != d.Declared || !sameFailure(before, d)
|
||||
}
|
||||
var since *time.Time
|
||||
times := 0
|
||||
if d.Outcome != OutcomeApplied {
|
||||
var before Doing
|
||||
var beforeFailed []byte
|
||||
err := i.store.Pool().QueryRow(ctx,
|
||||
`select outcome, refused, failed, failing_since, failures from node_report where node = $1`,
|
||||
node).Scan(&before.Outcome, &before.Refused, &beforeFailed, &before.Since, &before.Times)
|
||||
switch {
|
||||
case errors.Is(err, pgx.ErrNoRows):
|
||||
case err != nil:
|
||||
return err
|
||||
default:
|
||||
if err := json.Unmarshal(beforeFailed, &before.Failed); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
now := time.Now()
|
||||
since, times = &now, 1
|
||||
if err == nil && sameFailure(before, d) && before.Since != nil {
|
||||
if found == nil && sameFailure(before, d) && before.Since != nil {
|
||||
since, times = before.Since, before.Times+1
|
||||
}
|
||||
}
|
||||
@@ -707,7 +715,10 @@ func (i *Inventory) RecordDoing(ctx context.Context, node string, d Doing) error
|
||||
declared = excluded.declared,
|
||||
failing_since = excluded.failing_since, failures = excluded.failures`,
|
||||
node, d.Outcome, d.Refused, failed, d.Applied, d.Declared, since, times)
|
||||
return err
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
return news, nil
|
||||
}
|
||||
|
||||
// NotDoingWhatTheyWereTold is every machine whose last report was not a clean apply.
|
||||
|
||||
@@ -30,6 +30,11 @@ type Bus interface {
|
||||
// (design 29 §4, the *state* shape).
|
||||
PublishDeclaration(ctx context.Context, node string, body []byte) error
|
||||
|
||||
// PublishSeatEvent states a fact under a role's own name, for the holder of that role. A
|
||||
// module's event is addressed to the module; a role's is addressed to the role, so it keeps
|
||||
// meaning when the holder changes (novox/hq ADR 0121, ADR 0129).
|
||||
PublishSeatEvent(ctx context.Context, seat, event string, body []byte) error
|
||||
|
||||
// AskTool sends one question to a module's tool and awaits one answer. A tool nobody serves
|
||||
// must say so **at once** rather than after the whole wait: the difference between "that
|
||||
// module is down" and "that tool is slow" is the first thing a person asking wants.
|
||||
@@ -81,6 +86,13 @@ func EventSubject(source, key string) string {
|
||||
return "mesh.mod." + source + ".event." + key
|
||||
}
|
||||
|
||||
// SeatEventSubject is where a role's own event lands. Derived from the role, never from its holder:
|
||||
// a fact about the build machine or about the control plane keeps its address when the module holding
|
||||
// that role is replaced (novox/hq ADR 0121, ADR 0129).
|
||||
func SeatEventSubject(seat, event string) string {
|
||||
return "mesh.seat." + seat + ".event." + event
|
||||
}
|
||||
|
||||
// DeclareSubject is where one node's declaration lands. Last-per-subject on the NODES stream, so
|
||||
// a node that was away gets exactly the current one and a replayed older one is refused by
|
||||
// sequence — the wire-level answer to novox/hq issue 107.
|
||||
@@ -112,6 +124,27 @@ func (b OverNATS) PublishEvent(ctx context.Context, key, source, node string, bo
|
||||
return nil
|
||||
}
|
||||
|
||||
// PublishSeatEvent states a role's own fact. Same envelope as a module's event and a different
|
||||
// address: the source header is the role, because that is what the fact is about.
|
||||
func (b OverNATS) PublishSeatEvent(ctx context.Context, seat, event string, body []byte) error {
|
||||
id, err := eventID()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
h := nats.Header{}
|
||||
h.Set("x-event-id", id)
|
||||
h.Set("x-source", seat)
|
||||
_, err = b.JS.PublishMsg(&nats.Msg{
|
||||
Subject: SeatEventSubject(seat, event),
|
||||
Header: h,
|
||||
Data: body,
|
||||
}, nats.MsgId(id), nats.Context(ctx))
|
||||
if err != nil {
|
||||
return fmt.Errorf("stating %s of the %s seat: %w", event, seat, err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (b OverNATS) PublishDeclaration(ctx context.Context, node string, body []byte) error {
|
||||
_, err := b.JS.Publish(DeclareSubject(node), body, nats.Context(ctx))
|
||||
if err != nil {
|
||||
|
||||
+16
-13
@@ -265,7 +265,7 @@ func (e Enrolment) Outstanding(ctx context.Context, node string) (string, error)
|
||||
return e.Inventory.Outstanding(ctx, node)
|
||||
}
|
||||
|
||||
func (e Enrolment) Heard(ctx context.Context, report Report) (err error) {
|
||||
func (e Enrolment) Heard(ctx context.Context, report Report) (news bool, err error) {
|
||||
// A store that could not be asked right now is said as such, so the report is kept for
|
||||
// another attempt rather than acknowledged and lost (novox/hq issue 082).
|
||||
defer func() {
|
||||
@@ -274,11 +274,11 @@ func (e Enrolment) Heard(ctx context.Context, report Report) (err error) {
|
||||
}
|
||||
}()
|
||||
if report.Node == "" {
|
||||
return errors.New("a report named no node")
|
||||
return false, errors.New("a report named no node")
|
||||
}
|
||||
node, err := e.Inventory.NodeByName(ctx, report.Node)
|
||||
if err != nil {
|
||||
return err
|
||||
return false, err
|
||||
}
|
||||
|
||||
// What an adopted node holds, which firewall it found, and what is reachable on it (novox/hq
|
||||
@@ -298,7 +298,7 @@ func (e Enrolment) Heard(ctx context.Context, report Report) (err error) {
|
||||
Port: r.Port, By: r.By, Published: r.Published, ContainerPort: r.ContainerPort})
|
||||
}
|
||||
if err := e.Inventory.RecordAdoption(ctx, node.ID, held, report.Firewall, reachable); err != nil {
|
||||
return err
|
||||
return false, err
|
||||
}
|
||||
}
|
||||
// What it says about the tunnel it carried (novox/hq ADR 0105), whenever it says it.
|
||||
@@ -308,7 +308,7 @@ func (e Enrolment) Heard(ctx context.Context, report Report) (err error) {
|
||||
Peers: report.Tunnel.Peers, State: report.Tunnel.State, Note: report.Tunnel.Note,
|
||||
Kept: report.Tunnel.Kept,
|
||||
}); err != nil {
|
||||
return err
|
||||
return false, err
|
||||
}
|
||||
}
|
||||
// A node taking a found tunnel's key after enrolment (novox/hq ADR 0105). Verified against the
|
||||
@@ -317,9 +317,9 @@ func (e Enrolment) Heard(ctx context.Context, report Report) (err error) {
|
||||
// does not verify or is stale — a refusal, not "not now", so the node hears why.
|
||||
if report.Rekey != nil {
|
||||
if err := e.rekey(ctx, node, *report.Rekey); err != nil {
|
||||
return err
|
||||
return false, err
|
||||
}
|
||||
return e.Inventory.Seen(ctx, node.ID)
|
||||
return false, e.Inventory.Seen(ctx, node.ID)
|
||||
}
|
||||
|
||||
// A bare word that a node is there is not an account of what the machine did or holds: it
|
||||
@@ -334,7 +334,7 @@ func (e Enrolment) Heard(ctx context.Context, report Report) (err error) {
|
||||
if report.Superseded != "" {
|
||||
log.Printf("%s set aside declaration %s for the newer %s", report.Node, report.Declared, report.Superseded)
|
||||
}
|
||||
return e.Inventory.Seen(ctx, node.ID)
|
||||
return false, e.Inventory.Seen(ctx, node.ID)
|
||||
}
|
||||
// What it did is kept whichever way it went. Until this, a refusal or a failure moved
|
||||
// last_seen and the reason went to a log line, so "which machine is not doing what it was
|
||||
@@ -361,17 +361,20 @@ func (e Enrolment) Heard(ctx context.Context, report Report) (err error) {
|
||||
// on top of it (novox/hq ADR 0038). Kept even when the declaration was refused: what the
|
||||
// machine carries is true regardless of what it thought of the last thing it was sent.
|
||||
if err := e.Inventory.RecordCarried(ctx, report.Node, report.Carried); err != nil {
|
||||
return err
|
||||
return false, err
|
||||
}
|
||||
if err := e.Inventory.RecordDoing(ctx, node.ID, doing); err != nil {
|
||||
return err
|
||||
// **Whether this is news** is the store's answer: it holds the previous report, and a machine
|
||||
// that reconciles every minute says the same thing until something changes (novox/hq ADR 0134).
|
||||
news, err = e.Inventory.RecordDoing(ctx, node.ID, doing)
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
|
||||
// A refusal, a failure, or a bare word that the node is there — none of them is an account of
|
||||
// what the machine holds, so each moves last_seen and nothing else. Recording a partial list
|
||||
// as though it were the whole would tell a rebuilding node to remove what it still has.
|
||||
if report.Refused != "" || len(report.Failed) > 0 || report.Applied == nil {
|
||||
return e.Inventory.Seen(ctx, node.ID)
|
||||
return news, e.Inventory.Seen(ctx, node.ID)
|
||||
}
|
||||
return e.Inventory.RecordOwned(ctx, node.ID, report.Applied)
|
||||
return news, e.Inventory.RecordOwned(ctx, node.ID, report.Applied)
|
||||
}
|
||||
|
||||
@@ -49,6 +49,39 @@ func eventID() (string, error) {
|
||||
return hex.EncodeToString(raw), nil
|
||||
}
|
||||
|
||||
// MeshControllerSeat is the role the control plane holds, and therefore where its own facts live: a
|
||||
// role's events belong to the role, not to whichever container is holding it today (novox/hq ADR 0121,
|
||||
// ADR 0129). It is what makes them addressable while the control plane itself is being replaced.
|
||||
const MeshControllerSeat = "mesh-controller"
|
||||
|
||||
// The facts the mesh states about its own work (novox/hq ADR 0134).
|
||||
const (
|
||||
// KeyApplied: a machine now runs what it was sent.
|
||||
KeyApplied = "applied"
|
||||
// KeyRefused: a machine did not take what it was sent, and why.
|
||||
KeyRefused = "refused"
|
||||
// KeyBuiltBefore: a build the mesh already held, for a catalogue that asked what it missed. Not
|
||||
// `built` — that is the build machine's, said as it happens, and a replay is neither.
|
||||
KeyBuiltBefore = "built-before"
|
||||
)
|
||||
|
||||
// Applied is what a machine now runs, as the mesh states it.
|
||||
type Applied struct {
|
||||
Node string `json:"node"`
|
||||
Declared string `json:"declared,omitempty"`
|
||||
// Resources is how many the machine applied, not which: the list is the machine's own account
|
||||
// of itself and belongs in the records, not in a fact every listener has to read past.
|
||||
Resources int `json:"resources"`
|
||||
}
|
||||
|
||||
// Refused is a machine that would not take what it was sent.
|
||||
type Refused struct {
|
||||
Node string `json:"node"`
|
||||
Declared string `json:"declared,omitempty"`
|
||||
Refused string `json:"refused,omitempty"`
|
||||
Failed map[string]string `json:"failed,omitempty"`
|
||||
}
|
||||
|
||||
// KeyModuleBuilt is what the builder announces when it has built something. The catalogue places
|
||||
// it in the module graph; nothing else need care.
|
||||
const KeyModuleBuilt = "module.builder.built"
|
||||
|
||||
@@ -22,7 +22,7 @@ func heardFrom(t *testing.T, report link.Report) (*inventory.Inventory, inventor
|
||||
if _, err := inv.AddNode(ctx, report.Node); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := (link.Enrolment{Inventory: inv}).Heard(ctx, report); err != nil {
|
||||
if _, err := (link.Enrolment{Inventory: inv}).Heard(ctx, report); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
doing, said, err := inv.DoingOf(ctx, report.Node)
|
||||
@@ -103,7 +103,7 @@ func TestABareAliveDoesNotWipeTheDeclarationThatSaysANodeIsCurrent(t *testing.T)
|
||||
if err := inv.RecordSent(ctx, node.ID, digest); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := (link.Enrolment{Inventory: inv}).Heard(ctx, link.Report{
|
||||
if _, err := (link.Enrolment{Inventory: inv}).Heard(ctx, link.Report{
|
||||
Node: "anchor", Applied: []string{"a", "b"}, Declared: digest, Carried: []int{5432},
|
||||
}); err != nil {
|
||||
t.Fatal(err)
|
||||
@@ -126,7 +126,7 @@ func TestABareAliveDoesNotWipeTheDeclarationThatSaysANodeIsCurrent(t *testing.T)
|
||||
}
|
||||
|
||||
// Now the node says only that it is there, as it does every minute.
|
||||
if err := (link.Enrolment{Inventory: inv}).Heard(ctx, link.Report{Node: "anchor"}); err != nil {
|
||||
if _, err := (link.Enrolment{Inventory: inv}).Heard(ctx, link.Report{Node: "anchor"}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if !currentOf("anchor") {
|
||||
@@ -162,7 +162,7 @@ func TestAFailureDoesNotBecomeTheAccountOfWhatTheMachineHolds(t *testing.T) {
|
||||
if err := inv.RecordOwned(ctx, node.ID, []string{"one", "two", "three"}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := (link.Enrolment{Inventory: inv}).Heard(ctx, link.Report{
|
||||
if _, err := (link.Enrolment{Inventory: inv}).Heard(ctx, link.Report{
|
||||
Node: "workstation", Applied: []string{"one"}, Failed: map[string]string{"two": "no"},
|
||||
}); err != nil {
|
||||
t.Fatal(err)
|
||||
@@ -200,13 +200,13 @@ func TestWhatAnAdoptedNodeHoldsIsKeptAndAnAliveWordDoesNotWipeIt(t *testing.T) {
|
||||
}
|
||||
check("after the report")
|
||||
|
||||
if err := (link.Enrolment{Inventory: inv}).Heard(ctx, link.Report{Node: "anchor"}); err != nil {
|
||||
if _, err := (link.Enrolment{Inventory: inv}).Heard(ctx, link.Report{Node: "anchor"}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
check("after an alive word")
|
||||
|
||||
// A reconcile report carrying only adoption is recorded, though it applied nothing.
|
||||
if err := (link.Enrolment{Inventory: inv}).Heard(ctx, link.Report{Node: "anchor",
|
||||
if _, err := (link.Enrolment{Inventory: inv}).Heard(ctx, link.Report{Node: "anchor",
|
||||
Firewall: "ufw"}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
@@ -86,14 +86,15 @@ type counted struct {
|
||||
heard []Report
|
||||
}
|
||||
|
||||
func (c *counted) Heard(_ context.Context, r Report) error {
|
||||
func (c *counted) Heard(_ context.Context, r Report) (bool, error) {
|
||||
c.mu.Lock()
|
||||
defer c.mu.Unlock()
|
||||
if c.err != nil {
|
||||
return c.err
|
||||
return false, c.err
|
||||
}
|
||||
c.heard = append(c.heard, r)
|
||||
return nil
|
||||
// News, so what the mesh states about a report is exercised wherever a report is.
|
||||
return true, nil
|
||||
}
|
||||
|
||||
func (c *counted) refusing(err error) {
|
||||
@@ -207,11 +208,11 @@ type sentAndHeardSafely struct {
|
||||
heard []Report
|
||||
}
|
||||
|
||||
func (s *sentAndHeardSafely) Heard(_ context.Context, r Report) error {
|
||||
func (s *sentAndHeardSafely) Heard(_ context.Context, r Report) (bool, error) {
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
s.heard = append(s.heard, r)
|
||||
return nil
|
||||
return true, nil
|
||||
}
|
||||
|
||||
func (s *sentAndHeardSafely) Outstanding(context.Context, string) (string, error) {
|
||||
@@ -451,4 +452,108 @@ func TestNatsWorkSlowerThanTheWindowIsNotHandedOverAgain(t *testing.T) {
|
||||
// slowly is a listener that runs whatever it was given.
|
||||
type slowly struct{ work func() }
|
||||
|
||||
func (s slowly) Heard(context.Context, Report) error { s.work(); return nil }
|
||||
func (s slowly) Heard(context.Context, Report) (bool, error) { s.work(); return true, nil }
|
||||
|
||||
// **The mesh says what it applied** (novox/hq ADR 0134), under the seat the control plane holds — and
|
||||
// says nothing when a report is the same state said again, which is what a machine reconciling every
|
||||
// minute sends.
|
||||
func TestNatsTheMeshSaysWhatAMachineApplied(t *testing.T) {
|
||||
js := aBus(t)
|
||||
heard := make(chan *nats.Msg, 4)
|
||||
sub, err := js.Conn().Subscribe(SeatEventSubject(MeshControllerSeat, ">"), func(m *nats.Msg) {
|
||||
heard <- m
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer sub.Unsubscribe() //nolint:errcheck // the subscription dies with the connection
|
||||
|
||||
_, stop := servingOn(t, js, &counted{})
|
||||
defer stop()
|
||||
|
||||
// A report that changed something: the store says it was news.
|
||||
body, err := json.Marshal(Report{Node: "anchor", Declared: "d1", Applied: []string{"store", "broker"}})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := js.Context().Publish(ReportSubject("anchor"), body); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
select {
|
||||
case m := <-heard:
|
||||
if m.Subject != SeatEventSubject(MeshControllerSeat, KeyApplied) {
|
||||
t.Fatalf("the mesh stated %q", m.Subject)
|
||||
}
|
||||
var said Applied
|
||||
if err := json.Unmarshal(m.Data, &said); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if said.Node != "anchor" || said.Declared != "d1" || said.Resources != 2 {
|
||||
t.Fatalf("it said %+v", said)
|
||||
}
|
||||
case <-time.After(10 * time.Second):
|
||||
t.Fatal("the mesh said nothing about a machine that now runs something else")
|
||||
}
|
||||
|
||||
// A refusal is its own fact, with the reason in it rather than only in a log.
|
||||
refusal, err := json.Marshal(Report{Node: "anchor", Declared: "d2",
|
||||
Failed: map[string]string{"gitea.server": "no such image"}})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := js.Context().Publish(ReportSubject("anchor"), refusal); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
select {
|
||||
case m := <-heard:
|
||||
if m.Subject != SeatEventSubject(MeshControllerSeat, KeyRefused) {
|
||||
t.Fatalf("a refusal was stated as %q", m.Subject)
|
||||
}
|
||||
var said Refused
|
||||
if err := json.Unmarshal(m.Data, &said); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if said.Failed["gitea.server"] == "" {
|
||||
t.Fatalf("the refusal does not say which resource or why: %+v", said)
|
||||
}
|
||||
case <-time.After(10 * time.Second):
|
||||
t.Fatal("the mesh said nothing about a machine that refused what it was sent")
|
||||
}
|
||||
}
|
||||
|
||||
// And a report that is not news is not a fact. A machine reconciles every minute; a fact per report
|
||||
// would be a fact per minute per machine, which is a stream nobody reads.
|
||||
func TestNatsAReportThatIsNotNewsIsNotStated(t *testing.T) {
|
||||
js := aBus(t)
|
||||
heard := make(chan *nats.Msg, 4)
|
||||
sub, err := js.Conn().Subscribe(SeatEventSubject(MeshControllerSeat, ">"), func(m *nats.Msg) {
|
||||
heard <- m
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer sub.Unsubscribe() //nolint:errcheck // the subscription dies with the connection
|
||||
|
||||
// A store that records the report and says it was nothing new — which is what the mesh's own
|
||||
// store says about a machine repeating itself.
|
||||
_, stop := servingOn(t, js, sameAgain{})
|
||||
defer stop()
|
||||
|
||||
body, err := json.Marshal(Report{Node: "anchor", Declared: "d1", Applied: []string{"store"}})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := js.Context().Publish(ReportSubject("anchor"), body); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
select {
|
||||
case m := <-heard:
|
||||
t.Fatalf("the mesh stated %q about a machine that changed nothing", m.Subject)
|
||||
case <-time.After(3 * time.Second):
|
||||
}
|
||||
}
|
||||
|
||||
// sameAgain records a report and says it was the same state said again.
|
||||
type sameAgain struct{}
|
||||
|
||||
func (sameAgain) Heard(context.Context, Report) (bool, error) { return false, nil }
|
||||
|
||||
@@ -60,7 +60,7 @@ func TestASignedRekeyMovesTheHubOntoItsTunnel(t *testing.T) {
|
||||
rekey := &link.Rekey{Previous: ownKey, OverlayKey: tunnelKey, Tunnel: theTunnel()}
|
||||
rekey.Proof = ed25519.Sign(private, link.RekeyProof("anchor", ownKey, tunnelKey, theTunnel()))
|
||||
|
||||
if err := e.Heard(ctx, link.Report{Node: "anchor", Rekey: rekey}); err != nil {
|
||||
if _, err := e.Heard(ctx, link.Report{Node: "anchor", Rekey: rekey}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
placed, err := e.Inventory.Overlays(ctx)
|
||||
@@ -77,7 +77,7 @@ func TestASignedRekeyMovesTheHubOntoItsTunnel(t *testing.T) {
|
||||
_ = hub
|
||||
|
||||
// Replayed, it is stale: the previous key it names is no longer the node's.
|
||||
err = e.Heard(ctx, link.Report{Node: "anchor", Rekey: rekey})
|
||||
_, err = e.Heard(ctx, link.Report{Node: "anchor", Rekey: rekey})
|
||||
if err == nil || !strings.Contains(err.Error(), "previous overlay key") {
|
||||
t.Fatalf("a replayed rekey was accepted: %v", err)
|
||||
}
|
||||
@@ -93,7 +93,7 @@ func TestARekeySignedByAnotherKeyIsRefusedAndChangesNothing(t *testing.T) {
|
||||
rekey := &link.Rekey{Previous: ownKey, OverlayKey: tunnelKey, Tunnel: theTunnel()}
|
||||
rekey.Proof = ed25519.Sign(stranger, link.RekeyProof("anchor", ownKey, tunnelKey, theTunnel()))
|
||||
|
||||
err = e.Heard(ctx, link.Report{Node: "anchor", Rekey: rekey})
|
||||
_, err = e.Heard(ctx, link.Report{Node: "anchor", Rekey: rekey})
|
||||
if err == nil || !strings.Contains(err.Error(), "not signed by anchor's identity key") {
|
||||
t.Fatalf("a rekey signed by a stranger was accepted: %v", err)
|
||||
}
|
||||
@@ -111,7 +111,7 @@ func TestARekeySignedByAnotherKeyIsRefusedAndChangesNothing(t *testing.T) {
|
||||
other := theTunnel()
|
||||
other.Port = 51820
|
||||
moved.Proof = ed25519.Sign(mustPrivate(t, e, "anchor"), link.RekeyProof("anchor", ownKey, tunnelKey, other))
|
||||
if err := e.Heard(ctx, link.Report{Node: "anchor", Rekey: moved}); err == nil {
|
||||
if _, err := e.Heard(ctx, link.Report{Node: "anchor", Rekey: moved}); err == nil {
|
||||
t.Fatal("a proof over another tunnel was accepted")
|
||||
}
|
||||
}
|
||||
|
||||
@@ -9,12 +9,12 @@ import (
|
||||
|
||||
type heardWith struct{ err error }
|
||||
|
||||
func (h heardWith) Heard(context.Context, Report) error { return h.err }
|
||||
func (h heardWith) Heard(context.Context, Report) (bool, error) { return h.err == nil, h.err }
|
||||
|
||||
// switchable answers with whatever it is set to — the store away, then back.
|
||||
type switchable struct{ err error }
|
||||
|
||||
func (h *switchable) Heard(context.Context, Report) error { return h.err }
|
||||
func (h *switchable) Heard(context.Context, Report) (bool, error) { return h.err == nil, h.err }
|
||||
|
||||
func aReport(node, declared string) Report {
|
||||
return Report{Node: node, Declared: declared, Applied: []string{"store"}}
|
||||
|
||||
+62
-4
@@ -29,7 +29,12 @@ type Enroller interface {
|
||||
// Listener is what the controller does with a report. Separate from Enroller so the two can be
|
||||
// given independently, and so a server that only sends declarations needs neither.
|
||||
type Listener interface {
|
||||
Heard(ctx context.Context, report Report) error
|
||||
// Heard records what a node said, and says whether it was **news** — a machine that now runs
|
||||
// something else, or refuses something it did not refuse before. A machine reconciles
|
||||
// continuously and reports each time, so what is news is the store's answer rather than the
|
||||
// bus's: only this side has the previous report to compare with. What the mesh states about it
|
||||
// is the server's (novox/hq ADR 0134).
|
||||
Heard(ctx context.Context, report Report) (news bool, err error)
|
||||
}
|
||||
|
||||
// Recorder keeps what builders say.
|
||||
@@ -257,7 +262,7 @@ func (s *Server) heartbeat(m Control) {
|
||||
return
|
||||
}
|
||||
if s.listener != nil {
|
||||
if err := s.listener.Heard(context.Background(), Report{Node: alive.Node}); err != nil {
|
||||
if _, err := s.listener.Heard(context.Background(), Report{Node: alive.Node}); err != nil {
|
||||
s.log.Printf("could not record that %s is here: %v", alive.Node, err)
|
||||
}
|
||||
}
|
||||
@@ -291,7 +296,7 @@ func (s *Server) reported(ctx context.Context, m Control) {
|
||||
return
|
||||
}
|
||||
|
||||
err := s.listener.Heard(context.Background(), report)
|
||||
news, err := s.listener.Heard(context.Background(), report)
|
||||
switch s.decide(ctx, m, what, declaredIn, outstanding, err) {
|
||||
case Hold:
|
||||
// Held, not settled, while the store cannot take it: the node reports an apply once,
|
||||
@@ -307,6 +312,13 @@ func (s *Server) reported(ctx context.Context, m Control) {
|
||||
// node whose recovery copy is silently older than it looks.
|
||||
s.log.Printf("could not record %s's report: %v", report.Node, err)
|
||||
}
|
||||
// **And the mesh says what it did** (novox/hq ADR 0134). Only when the report was news: a
|
||||
// machine reports every convergence, and a fact per report would be a fact per minute per
|
||||
// machine saying nothing. Stated after it is recorded, so nothing is announced that the
|
||||
// mesh does not hold.
|
||||
if err == nil && news {
|
||||
s.saysWhatItDid(ctx, report)
|
||||
}
|
||||
}
|
||||
|
||||
switch {
|
||||
@@ -460,7 +472,19 @@ func (s *Server) catchingUp(ctx context.Context, m Control) {
|
||||
sent := 0
|
||||
for _, a := range announcements {
|
||||
a.Replay = true
|
||||
if err := EmitEvent(ctx, s.bus, KeyModuleBuilt, "control-plane", "", a); err != nil {
|
||||
// Under the control plane's own seat (novox/hq ADR 0134). It used to be published as a
|
||||
// module's event from a module called "control-plane", which does not exist — so the
|
||||
// controller's own account refused it, every catalogue that asked what it missed was
|
||||
// answered with nothing, and its graph kept the gap (found 2026-09-28).
|
||||
body, err := json.Marshal(a)
|
||||
if err != nil {
|
||||
// A body that cannot be written is this program's fault, not the bus's, and publishing
|
||||
// an empty one would put a fact on the mesh that says nothing.
|
||||
s.log.Printf("cannot re-announce %s at %s: %v", a.Module, short(a.Commit), err)
|
||||
_ = m.Took()
|
||||
return
|
||||
}
|
||||
if err := s.bus.PublishSeatEvent(ctx, MeshControllerSeat, KeyBuiltBefore, body); err != nil {
|
||||
// Said and abandoned rather than retried: the catalogue asks again every time it
|
||||
// starts, and half a graph delivered twice is no better than half delivered once.
|
||||
s.log.Printf("replaying %s at %s failed, and the rest is abandoned: %v",
|
||||
@@ -551,3 +575,37 @@ func (s *Server) sourceMoved(ctx context.Context, m Control) {
|
||||
}
|
||||
_ = m.Took()
|
||||
}
|
||||
|
||||
// saysWhatItDid states what a machine now runs, or what it would not take, as a fact on the bus
|
||||
// (novox/hq ADR 0134).
|
||||
//
|
||||
// **The control plane speaks, as the holder of its seat.** A node's report is control traffic only
|
||||
// this process may read, so the chain from a merge to a machine went dark exactly where it touched
|
||||
// one: nothing said which version a machine runs, or that it refused to. The facts are second-hand
|
||||
// on purpose — one emitter, one ordering — and a machine that cannot reach the bus produces none, so
|
||||
// absence is not health.
|
||||
//
|
||||
// A failure to state a fact is logged and nothing else: the report is recorded, which is the part
|
||||
// that must not be lost, and the next change says the same thing again.
|
||||
func (s *Server) saysWhatItDid(ctx context.Context, report Report) {
|
||||
if s.bus == nil {
|
||||
return
|
||||
}
|
||||
event, body := KeyApplied, any(Applied{
|
||||
Node: report.Node, Declared: report.Declared, Resources: len(report.Applied),
|
||||
})
|
||||
if report.Refused != "" || len(report.Failed) > 0 {
|
||||
event, body = KeyRefused, Refused{
|
||||
Node: report.Node, Declared: report.Declared,
|
||||
Refused: report.Refused, Failed: report.Failed,
|
||||
}
|
||||
}
|
||||
raw, err := json.Marshal(body)
|
||||
if err != nil {
|
||||
s.log.Printf("could not say what %s did: %v", report.Node, err)
|
||||
return
|
||||
}
|
||||
if err := s.bus.PublishSeatEvent(ctx, MeshControllerSeat, event, raw); err != nil {
|
||||
s.log.Printf("could not say that %s %s: %v", report.Node, event, err)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -23,12 +23,12 @@ type sentAndHeard struct {
|
||||
err error
|
||||
}
|
||||
|
||||
func (s *sentAndHeard) Heard(_ context.Context, r Report) error {
|
||||
func (s *sentAndHeard) Heard(_ context.Context, r Report) (bool, error) {
|
||||
if s.err != nil {
|
||||
return s.err
|
||||
return false, s.err
|
||||
}
|
||||
s.heard = append(s.heard, r)
|
||||
return nil
|
||||
return true, nil
|
||||
}
|
||||
|
||||
func (s *sentAndHeard) Outstanding(context.Context, string) (string, error) { return s.sent, nil }
|
||||
|
||||
Reference in New Issue
Block a user