Compare commits
9
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
ce9e20fbbc | ||
|
|
878690697e | ||
|
|
ad97297576 | ||
|
|
683b1ed693 | ||
|
|
04f9f378b0 | ||
|
|
1c3f44a526 | ||
|
|
89e152dfe2 | ||
|
|
1ebad3786c | ||
|
|
f2f526a60a |
@@ -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 {
|
if err := open.inventory.RecordSent(ctx, record.ID, digestOf(body)); err != nil {
|
||||||
t.Fatal(err)
|
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),
|
Node: "anchor", Applied: []string{"hello-web.x"}, Declared: digestOf(body),
|
||||||
Firewall: "ufw", Held: held, Reachable: reachable,
|
Firewall: "ufw", Held: held, Reachable: reachable,
|
||||||
}); err != nil {
|
}); err != nil {
|
||||||
|
|||||||
@@ -162,7 +162,13 @@ func declare(ctx context.Context, args []string) error {
|
|||||||
return err
|
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 {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -181,6 +181,14 @@ func PermissionsFor(p Principal) (Permissions, error) {
|
|||||||
for _, seat := range meshSeatsTheControllerUses {
|
for _, seat := range meshSeatsTheControllerUses {
|
||||||
pub = append(pub, "mesh.seat."+seat+".accept.>")
|
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
|
// 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
|
// 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
|
// 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.>`).
|
// ack subject is derived from (nats.go: `$JS.ACK.<stream>.controller.>`).
|
||||||
const ControllerName = "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
|
// 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.
|
// 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
|
jetstream: enabled
|
||||||
users = [
|
users = [
|
||||||
{ user: "controller", password: "$2a$11$cccccccccccccccccccccc", permissions: {
|
{ 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"] }
|
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" }
|
allow_responses: { max: 1, ttl: "1m" }
|
||||||
} }
|
} }
|
||||||
|
|||||||
@@ -137,7 +137,7 @@ func (r Registry) has(ctx context.Context, url string, accept ...string) (bool,
|
|||||||
return false, err
|
return false, err
|
||||||
}
|
}
|
||||||
for _, media := range accept {
|
for _, media := range accept {
|
||||||
request.Header.Set("Accept", media)
|
request.Header.Add("Accept", media)
|
||||||
}
|
}
|
||||||
response, err := r.client().Do(request)
|
response, err := r.client().Do(request)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
|||||||
@@ -48,7 +48,11 @@ type Seat struct {
|
|||||||
//
|
//
|
||||||
// In the order a person reads it: the mesh's own, then a node's.
|
// In the order a person reads it: the mesh's own, then a node's.
|
||||||
var defaultSeats = []Seat{
|
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"},
|
{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
|
// **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
|
// 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")
|
refuser := nodeNamed(t, inv, "refuser")
|
||||||
failer := nodeNamed(t, inv, "failer")
|
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",
|
Outcome: OutcomeRefused, Refused: "resource \"x\": a file needs a path",
|
||||||
}); err != nil {
|
}); err != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
if err := inv.RecordDoing(ctx, failer, Doing{
|
if _, err := inv.RecordDoing(ctx, failer, Doing{
|
||||||
Outcome: OutcomeFailed,
|
Outcome: OutcomeFailed,
|
||||||
Failed: []FailedResource{{ID: "svc", Error: "unit not found"}},
|
Failed: []FailedResource{{ID: "svc", Error: "unit not found"}},
|
||||||
Applied: 4,
|
Applied: 4,
|
||||||
@@ -70,7 +70,7 @@ func TestAMachineDoingWhatItWasToldIsNotOnTheList(t *testing.T) {
|
|||||||
inv := fresh(t)
|
inv := fresh(t)
|
||||||
ctx := context.Background()
|
ctx := context.Background()
|
||||||
id := nodeNamed(t, inv, "fine")
|
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)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
wrong, err := inv.NotDoingWhatTheyWereTold(ctx)
|
wrong, err := inv.NotDoingWhatTheyWereTold(ctx)
|
||||||
@@ -97,12 +97,12 @@ func TestTheLastReportReplacesTheOneBefore(t *testing.T) {
|
|||||||
inv := fresh(t)
|
inv := fresh(t)
|
||||||
ctx := context.Background()
|
ctx := context.Background()
|
||||||
id := nodeNamed(t, inv, "recovered")
|
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"}},
|
Outcome: OutcomeFailed, Failed: []FailedResource{{ID: "a", Error: "no"}},
|
||||||
}); err != nil {
|
}); err != nil {
|
||||||
t.Fatal(err)
|
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)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
wrong, err := inv.NotDoingWhatTheyWereTold(ctx)
|
wrong, err := inv.NotDoingWhatTheyWereTold(ctx)
|
||||||
@@ -141,7 +141,7 @@ func TestWhatANodeSaidGoesWhenTheNodeDoes(t *testing.T) {
|
|||||||
inv := fresh(t)
|
inv := fresh(t)
|
||||||
ctx := context.Background()
|
ctx := context.Background()
|
||||||
id := nodeNamed(t, inv, "leaving")
|
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)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
if _, err := inv.store.Pool().Exec(ctx, `delete from node where name = 'leaving'`); err != nil {
|
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")
|
id := nodeNamed(t, inv, "looping")
|
||||||
same := Doing{Outcome: OutcomeFailed, Failed: []FailedResource{{ID: "img", Error: "no such image"}}}
|
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)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
first, _, err := inv.DoingOf(ctx, "looping")
|
first, _, err := inv.DoingOf(ctx, "looping")
|
||||||
@@ -269,7 +269,7 @@ func TestTheSameFailureReportedAgainIsCountedNotRestarted(t *testing.T) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
for range StuckAfter - 1 {
|
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)
|
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
|
// 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.
|
// failure: it is the resource that loops, not the sentence.
|
||||||
reworded := Doing{Outcome: OutcomeFailed, Failed: []FailedResource{{ID: "img", Error: "no such image (after 31s)"}}}
|
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)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
still, _, err := inv.DoingOf(ctx, "looping")
|
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.
|
// A different failure is a new situation, not a longer one.
|
||||||
other := Doing{Outcome: OutcomeFailed, Failed: []FailedResource{{ID: "svc", Error: "unit not found"}}}
|
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)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
changed, _, err := inv.DoingOf(ctx, "looping")
|
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.
|
// 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)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
fine, _, err := inv.DoingOf(ctx, "looping")
|
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.
|
// 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)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
wrong, err := inv.NotDoingWhatTheyWereTold(ctx)
|
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
|
// 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
|
// 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.
|
// 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)
|
failed, err := json.Marshal(d.Failed)
|
||||||
if err != nil {
|
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
|
var since *time.Time
|
||||||
times := 0
|
times := 0
|
||||||
if d.Outcome != OutcomeApplied {
|
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()
|
now := time.Now()
|
||||||
since, times = &now, 1
|
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
|
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,
|
declared = excluded.declared,
|
||||||
failing_since = excluded.failing_since, failures = excluded.failures`,
|
failing_since = excluded.failing_since, failures = excluded.failures`,
|
||||||
node, d.Outcome, d.Refused, failed, d.Applied, d.Declared, since, times)
|
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.
|
// 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).
|
// (design 29 §4, the *state* shape).
|
||||||
PublishDeclaration(ctx context.Context, node string, body []byte) error
|
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
|
// 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
|
// 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.
|
// 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
|
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
|
// 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
|
// 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.
|
// 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
|
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 {
|
func (b OverNATS) PublishDeclaration(ctx context.Context, node string, body []byte) error {
|
||||||
_, err := b.JS.Publish(DeclareSubject(node), body, nats.Context(ctx))
|
_, err := b.JS.Publish(DeclareSubject(node), body, nats.Context(ctx))
|
||||||
if err != nil {
|
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)
|
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
|
// 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).
|
// another attempt rather than acknowledged and lost (novox/hq issue 082).
|
||||||
defer func() {
|
defer func() {
|
||||||
@@ -274,11 +274,11 @@ func (e Enrolment) Heard(ctx context.Context, report Report) (err error) {
|
|||||||
}
|
}
|
||||||
}()
|
}()
|
||||||
if report.Node == "" {
|
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)
|
node, err := e.Inventory.NodeByName(ctx, report.Node)
|
||||||
if err != nil {
|
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
|
// 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})
|
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 {
|
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.
|
// 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,
|
Peers: report.Tunnel.Peers, State: report.Tunnel.State, Note: report.Tunnel.Note,
|
||||||
Kept: report.Tunnel.Kept,
|
Kept: report.Tunnel.Kept,
|
||||||
}); err != nil {
|
}); err != nil {
|
||||||
return err
|
return false, err
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
// A node taking a found tunnel's key after enrolment (novox/hq ADR 0105). Verified against the
|
// 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.
|
// does not verify or is stale — a refusal, not "not now", so the node hears why.
|
||||||
if report.Rekey != nil {
|
if report.Rekey != nil {
|
||||||
if err := e.rekey(ctx, node, *report.Rekey); err != 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
|
// 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 != "" {
|
if report.Superseded != "" {
|
||||||
log.Printf("%s set aside declaration %s for the newer %s", report.Node, report.Declared, 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
|
// 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
|
// 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
|
// 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.
|
// 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 {
|
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 {
|
// **Whether this is news** is the store's answer: it holds the previous report, and a machine
|
||||||
return err
|
// 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
|
// 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
|
// 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.
|
// 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 {
|
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
|
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
|
// KeyModuleBuilt is what the builder announces when it has built something. The catalogue places
|
||||||
// it in the module graph; nothing else need care.
|
// it in the module graph; nothing else need care.
|
||||||
const KeyModuleBuilt = "module.builder.built"
|
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 {
|
if _, err := inv.AddNode(ctx, report.Node); err != nil {
|
||||||
t.Fatal(err)
|
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)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
doing, said, err := inv.DoingOf(ctx, report.Node)
|
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 {
|
if err := inv.RecordSent(ctx, node.ID, digest); err != nil {
|
||||||
t.Fatal(err)
|
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},
|
Node: "anchor", Applied: []string{"a", "b"}, Declared: digest, Carried: []int{5432},
|
||||||
}); err != nil {
|
}); err != nil {
|
||||||
t.Fatal(err)
|
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.
|
// 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)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
if !currentOf("anchor") {
|
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 {
|
if err := inv.RecordOwned(ctx, node.ID, []string{"one", "two", "three"}); err != nil {
|
||||||
t.Fatal(err)
|
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"},
|
Node: "workstation", Applied: []string{"one"}, Failed: map[string]string{"two": "no"},
|
||||||
}); err != nil {
|
}); err != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
@@ -200,13 +200,13 @@ func TestWhatAnAdoptedNodeHoldsIsKeptAndAnAliveWordDoesNotWipeIt(t *testing.T) {
|
|||||||
}
|
}
|
||||||
check("after the report")
|
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)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
check("after an alive word")
|
check("after an alive word")
|
||||||
|
|
||||||
// A reconcile report carrying only adoption is recorded, though it applied nothing.
|
// 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 {
|
Firewall: "ufw"}); err != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -86,14 +86,15 @@ type counted struct {
|
|||||||
heard []Report
|
heard []Report
|
||||||
}
|
}
|
||||||
|
|
||||||
func (c *counted) Heard(_ context.Context, r Report) error {
|
func (c *counted) Heard(_ context.Context, r Report) (bool, error) {
|
||||||
c.mu.Lock()
|
c.mu.Lock()
|
||||||
defer c.mu.Unlock()
|
defer c.mu.Unlock()
|
||||||
if c.err != nil {
|
if c.err != nil {
|
||||||
return c.err
|
return false, c.err
|
||||||
}
|
}
|
||||||
c.heard = append(c.heard, r)
|
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) {
|
func (c *counted) refusing(err error) {
|
||||||
@@ -207,11 +208,11 @@ type sentAndHeardSafely struct {
|
|||||||
heard []Report
|
heard []Report
|
||||||
}
|
}
|
||||||
|
|
||||||
func (s *sentAndHeardSafely) Heard(_ context.Context, r Report) error {
|
func (s *sentAndHeardSafely) Heard(_ context.Context, r Report) (bool, error) {
|
||||||
s.mu.Lock()
|
s.mu.Lock()
|
||||||
defer s.mu.Unlock()
|
defer s.mu.Unlock()
|
||||||
s.heard = append(s.heard, r)
|
s.heard = append(s.heard, r)
|
||||||
return nil
|
return true, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func (s *sentAndHeardSafely) Outstanding(context.Context, string) (string, error) {
|
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.
|
// slowly is a listener that runs whatever it was given.
|
||||||
type slowly struct{ work func() }
|
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 := &link.Rekey{Previous: ownKey, OverlayKey: tunnelKey, Tunnel: theTunnel()}
|
||||||
rekey.Proof = ed25519.Sign(private, link.RekeyProof("anchor", ownKey, tunnelKey, 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)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
placed, err := e.Inventory.Overlays(ctx)
|
placed, err := e.Inventory.Overlays(ctx)
|
||||||
@@ -77,7 +77,7 @@ func TestASignedRekeyMovesTheHubOntoItsTunnel(t *testing.T) {
|
|||||||
_ = hub
|
_ = hub
|
||||||
|
|
||||||
// Replayed, it is stale: the previous key it names is no longer the node's.
|
// 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") {
|
if err == nil || !strings.Contains(err.Error(), "previous overlay key") {
|
||||||
t.Fatalf("a replayed rekey was accepted: %v", err)
|
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 := &link.Rekey{Previous: ownKey, OverlayKey: tunnelKey, Tunnel: theTunnel()}
|
||||||
rekey.Proof = ed25519.Sign(stranger, link.RekeyProof("anchor", ownKey, tunnelKey, 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") {
|
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)
|
t.Fatalf("a rekey signed by a stranger was accepted: %v", err)
|
||||||
}
|
}
|
||||||
@@ -111,7 +111,7 @@ func TestARekeySignedByAnotherKeyIsRefusedAndChangesNothing(t *testing.T) {
|
|||||||
other := theTunnel()
|
other := theTunnel()
|
||||||
other.Port = 51820
|
other.Port = 51820
|
||||||
moved.Proof = ed25519.Sign(mustPrivate(t, e, "anchor"), link.RekeyProof("anchor", ownKey, tunnelKey, other))
|
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")
|
t.Fatal("a proof over another tunnel was accepted")
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -9,12 +9,12 @@ import (
|
|||||||
|
|
||||||
type heardWith struct{ err error }
|
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.
|
// switchable answers with whatever it is set to — the store away, then back.
|
||||||
type switchable struct{ err error }
|
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 {
|
func aReport(node, declared string) Report {
|
||||||
return Report{Node: node, Declared: declared, Applied: []string{"store"}}
|
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
|
// 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.
|
// given independently, and so a server that only sends declarations needs neither.
|
||||||
type Listener interface {
|
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.
|
// Recorder keeps what builders say.
|
||||||
@@ -257,7 +262,7 @@ func (s *Server) heartbeat(m Control) {
|
|||||||
return
|
return
|
||||||
}
|
}
|
||||||
if s.listener != nil {
|
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)
|
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
|
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) {
|
switch s.decide(ctx, m, what, declaredIn, outstanding, err) {
|
||||||
case Hold:
|
case Hold:
|
||||||
// Held, not settled, while the store cannot take it: the node reports an apply once,
|
// 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.
|
// node whose recovery copy is silently older than it looks.
|
||||||
s.log.Printf("could not record %s's report: %v", report.Node, err)
|
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 {
|
switch {
|
||||||
@@ -460,7 +472,19 @@ func (s *Server) catchingUp(ctx context.Context, m Control) {
|
|||||||
sent := 0
|
sent := 0
|
||||||
for _, a := range announcements {
|
for _, a := range announcements {
|
||||||
a.Replay = true
|
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
|
// 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.
|
// 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",
|
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()
|
_ = 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
|
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 {
|
if s.err != nil {
|
||||||
return s.err
|
return false, s.err
|
||||||
}
|
}
|
||||||
s.heard = append(s.heard, r)
|
s.heard = append(s.heard, r)
|
||||||
return nil
|
return true, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func (s *sentAndHeard) Outstanding(context.Context, string) (string, error) { return s.sent, nil }
|
func (s *sentAndHeard) Outstanding(context.Context, string) (string, error) { return s.sent, nil }
|
||||||
|
|||||||
Reference in New Issue
Block a user