Compare commits
3
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
bf82d30865 | ||
|
|
c432b9b5f6 | ||
|
|
7bd61332cb |
@@ -0,0 +1,207 @@
|
|||||||
|
package main
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"errors"
|
||||||
|
"fmt"
|
||||||
|
"sort"
|
||||||
|
|
||||||
|
"github.com/nats-io/nats.go"
|
||||||
|
|
||||||
|
"github.com/novox/mesh-controller/internal/conditions"
|
||||||
|
)
|
||||||
|
|
||||||
|
// **A filling bucket or log is said before it is full** (novox/hq ADR 0297 §6, issue 501).
|
||||||
|
//
|
||||||
|
// A module's bucket and its log each have a cap on the bus, and when either is full the bus refuses
|
||||||
|
// the next write: the module stops doing what it writes for, and nothing said so before. On
|
||||||
|
// 2026-10-11 the issue tracker's bucket held a sixth of its cap after two days, growing toward a stop
|
||||||
|
// nobody would have heard coming. So the self-check reads every module's bucket and log, and one at
|
||||||
|
// fillRaiseAt of its cap or more raises a condition naming the module, the bucket or log, and how full
|
||||||
|
// it is.
|
||||||
|
//
|
||||||
|
// **Cleared by observation below fillClearBelow, not below fillRaiseAt**: a bucket or log hovering at
|
||||||
|
// its threshold would otherwise be raised and cleared on every run. Between the two, an open condition
|
||||||
|
// is kept and none is raised.
|
||||||
|
//
|
||||||
|
// **One that cannot be read is said, and the rest are still judged.** An unanswered question about one
|
||||||
|
// stream is no reason to know nothing of the others: it is a finding of its own, naming the bucket or
|
||||||
|
// log and why, and a fill condition already open for it is kept, because not knowing is not a pass.
|
||||||
|
|
||||||
|
// The fill thresholds, in percent of a bucket's or log's cap.
|
||||||
|
const (
|
||||||
|
fillRaiseAt = 75
|
||||||
|
fillClearBelow = 70
|
||||||
|
)
|
||||||
|
|
||||||
|
// The conditions the fill probe raises: one filling, and one that could not be read.
|
||||||
|
const (
|
||||||
|
kindBucketOrLogFilling = "bucket-or-log-filling"
|
||||||
|
kindBucketOrLogUnread = "bucket-or-log-unread"
|
||||||
|
)
|
||||||
|
|
||||||
|
// probeFillID is the probe's id in the registry.
|
||||||
|
const probeFillID = "D-fill"
|
||||||
|
|
||||||
|
// bucketOrLog is one module's bucket or log as the fill probe reads it.
|
||||||
|
type bucketOrLog struct {
|
||||||
|
Module string
|
||||||
|
// Kind is `bucket` or `log`; Name its local name; Stream the stream it is on the bus.
|
||||||
|
Kind, Name, Stream string
|
||||||
|
}
|
||||||
|
|
||||||
|
// filled is how full one bucket or log is, read from the bus.
|
||||||
|
type filled struct {
|
||||||
|
Bytes, Max uint64
|
||||||
|
}
|
||||||
|
|
||||||
|
// percent is how full, in whole percent, rounded down.
|
||||||
|
func (f filled) percent() uint64 {
|
||||||
|
if f.Max == 0 {
|
||||||
|
return 0
|
||||||
|
}
|
||||||
|
return f.Bytes * 100 / f.Max
|
||||||
|
}
|
||||||
|
|
||||||
|
// fillKey is where a bucket's or log's fill condition is kept.
|
||||||
|
func fillKey(s bucketOrLog) string { return conditions.Key(conditions.ScopeBus, s.Stream, "filling") }
|
||||||
|
|
||||||
|
// fillObservation is what one bucket's or log's fill says: a finding at fillRaiseAt or above, and,
|
||||||
|
// while its condition is open, at fillClearBelow or above too; nothing otherwise, which is what clears
|
||||||
|
// it. The operator's words are the kind's (plain_words.go); the summary and the evidence name the
|
||||||
|
// module, the bucket or log, and the fill.
|
||||||
|
func fillObservation(s bucketOrLog, f filled, open bool) (conditions.Observation, bool) {
|
||||||
|
if f.Max == 0 {
|
||||||
|
return conditions.Observation{}, false // one with no cap cannot fill
|
||||||
|
}
|
||||||
|
// Compared in bytes, not rounded percent: 74.9% is not 75%.
|
||||||
|
atRaise := f.Bytes*100 >= f.Max*fillRaiseAt
|
||||||
|
aboveClear := f.Bytes*100 >= f.Max*fillClearBelow
|
||||||
|
if !atRaise && !(open && aboveClear) {
|
||||||
|
return conditions.Observation{}, false
|
||||||
|
}
|
||||||
|
pct := f.percent()
|
||||||
|
return conditions.Observation{
|
||||||
|
Scope: conditions.ScopeBus, ID: s.Stream, Token: "filling", Kind: kindBucketOrLogFilling,
|
||||||
|
Severity: conditions.Warning,
|
||||||
|
Summary: fmt.Sprintf("%s's %s %s holds %s of %s, %d%%: at its cap the bus refuses the module's next write",
|
||||||
|
s.Module, s.Kind, s.Name, mibWords(f.Bytes), mibWords(f.Max), pct),
|
||||||
|
Said: fmt.Sprintf("module %s, %s %s (stream %s): %d of %d bytes, %d%%", s.Module, s.Kind, s.Name,
|
||||||
|
s.Stream, f.Bytes, f.Max, pct),
|
||||||
|
}, true
|
||||||
|
}
|
||||||
|
|
||||||
|
// unreadObservation says one bucket or log could not be read, and why. Asked over the network, so one
|
||||||
|
// look can be wrong: raised on the second look in a row (confirm.go).
|
||||||
|
func unreadObservation(s bucketOrLog, why error) conditions.Observation {
|
||||||
|
return conditions.Observation{
|
||||||
|
Scope: conditions.ScopeBus, ID: s.Stream, Token: "unread", Kind: kindBucketOrLogUnread,
|
||||||
|
Severity: conditions.Warning, Confirm: true,
|
||||||
|
Summary: fmt.Sprintf("%s's %s %s could not be read, so how full it is is not known", s.Module, s.Kind, s.Name),
|
||||||
|
Said: fmt.Sprintf("module %s, %s %s (stream %s): %v", s.Module, s.Kind, s.Name, s.Stream, why),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// keptFilling is an open fill condition said again for a bucket or log that could not be read this
|
||||||
|
// time: not knowing how full it is now is no reason to say it has room. Only ever made from a condition
|
||||||
|
// read back as open — its own summary and severity, never words made up for it.
|
||||||
|
func keptFilling(s bucketOrLog, open conditions.Condition, why error) conditions.Observation {
|
||||||
|
return conditions.Observation{
|
||||||
|
Scope: conditions.ScopeBus, ID: s.Stream, Token: "filling", Kind: kindBucketOrLogFilling,
|
||||||
|
Severity: open.Severity, Summary: open.Summary,
|
||||||
|
Said: fmt.Sprintf("module %s, %s %s (stream %s): not read this time (%v); kept open as it was",
|
||||||
|
s.Module, s.Kind, s.Name, s.Stream, why),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// mibWords is a size as a person reads it, in MiB with one decimal.
|
||||||
|
func mibWords(b uint64) string {
|
||||||
|
return fmt.Sprintf("%.1f MiB", float64(b)/(1024*1024))
|
||||||
|
}
|
||||||
|
|
||||||
|
// readFill is how full one bucket or log is on the bus; found false when it is not on the bus.
|
||||||
|
type readFill func(ctx context.Context, s bucketOrLog) (f filled, found bool, err error)
|
||||||
|
|
||||||
|
// openFill is the fill condition open under a key, if one is, or why it could not be read.
|
||||||
|
type openFill func(ctx context.Context, key string) (conditions.Condition, bool, error)
|
||||||
|
|
||||||
|
// judgeFills judges every bucket and log on its own: one that cannot be read is said as unread, with
|
||||||
|
// its fill condition kept only when that condition was read back and is open — when it cannot be read
|
||||||
|
// either, nothing is said of its fill, which the unread finding already covers. Every other is judged
|
||||||
|
// by its fill; for one of those, a condition that cannot be read is taken as open, so an unknown never
|
||||||
|
// clears a fill measured between fillClearBelow and fillRaiseAt.
|
||||||
|
func judgeFills(ctx context.Context, all []bucketOrLog, read readFill, open openFill) []conditions.Observation {
|
||||||
|
var out []conditions.Observation
|
||||||
|
for _, s := range all {
|
||||||
|
f, found, err := read(ctx, s)
|
||||||
|
if err != nil {
|
||||||
|
out = append(out, unreadObservation(s, err))
|
||||||
|
if c, isOpen, cerr := open(ctx, fillKey(s)); cerr == nil && isOpen {
|
||||||
|
out = append(out, keptFilling(s, c, err))
|
||||||
|
}
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
if !found {
|
||||||
|
continue // not on the bus: created on the controller's next raise, and nothing there can fill
|
||||||
|
}
|
||||||
|
_, isOpen, cerr := open(ctx, fillKey(s))
|
||||||
|
isOpen = isOpen || cerr != nil
|
||||||
|
if o, said := fillObservation(s, f, isOpen); said {
|
||||||
|
out = append(out, o)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
sort.SliceStable(out, func(i, j int) bool { return out[i].Key() < out[j].Key() })
|
||||||
|
return out
|
||||||
|
}
|
||||||
|
|
||||||
|
// declaredBucketsAndLogs is every module's bucket and log the catalogue declares, by its stream.
|
||||||
|
func declaredBucketsAndLogs(ctx context.Context, d *doctor) ([]bucketOrLog, error) {
|
||||||
|
buckets, err := d.open.inventory.DeclaredBuckets(ctx)
|
||||||
|
if err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
logs, err := d.open.inventory.DeclaredLogs(ctx)
|
||||||
|
if err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
var out []bucketOrLog
|
||||||
|
for _, b := range buckets {
|
||||||
|
out = append(out, bucketOrLog{Module: b.Module, Kind: "bucket", Name: b.Name, Stream: "KV_" + b.Bucket()})
|
||||||
|
}
|
||||||
|
for _, l := range logs {
|
||||||
|
out = append(out, bucketOrLog{Module: l.Module, Kind: "log", Name: l.Name, Stream: l.Stream()})
|
||||||
|
}
|
||||||
|
return out, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// probeFill reads every module's bucket and log on the bus and says each one filling toward its cap,
|
||||||
|
// and each one that could not be read.
|
||||||
|
func probeFill(ctx context.Context, d *doctor) ([]conditions.Observation, error) {
|
||||||
|
all, err := declaredBucketsAndLogs(ctx, d)
|
||||||
|
if err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
js := d.js.Context()
|
||||||
|
read := func(ctx context.Context, s bucketOrLog) (filled, bool, error) {
|
||||||
|
if err := ctx.Err(); err != nil {
|
||||||
|
return filled{}, false, err
|
||||||
|
}
|
||||||
|
info, err := js.StreamInfo(s.Stream, nats.Context(ctx))
|
||||||
|
switch {
|
||||||
|
case errors.Is(err, nats.ErrStreamNotFound):
|
||||||
|
return filled{}, false, nil
|
||||||
|
case err != nil:
|
||||||
|
return filled{}, false, err
|
||||||
|
case info.Config.MaxBytes <= 0:
|
||||||
|
return filled{}, true, nil
|
||||||
|
}
|
||||||
|
return filled{Bytes: info.State.Bytes, Max: uint64(info.Config.MaxBytes)}, true, nil
|
||||||
|
}
|
||||||
|
open := func(ctx context.Context, key string) (conditions.Condition, bool, error) {
|
||||||
|
if d.keeper == nil {
|
||||||
|
return conditions.Condition{}, false, nil
|
||||||
|
}
|
||||||
|
return d.keeper.Get(ctx, key)
|
||||||
|
}
|
||||||
|
return judgeFills(ctx, all, read, open), nil
|
||||||
|
}
|
||||||
@@ -0,0 +1,165 @@
|
|||||||
|
package main
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"errors"
|
||||||
|
"strings"
|
||||||
|
"testing"
|
||||||
|
|
||||||
|
"github.com/novox/mesh-controller/internal/conditions"
|
||||||
|
)
|
||||||
|
|
||||||
|
// **A filling log or bucket is said at three quarters of its cap and cleared below seven tenths**
|
||||||
|
// (novox/hq ADR 0297 §6): raised at 75% naming the module, the bucket or log and its fill; not at 74.9%;
|
||||||
|
// kept between 70% and 75% while it is open, and not raised there when it is not; and cleared — said no
|
||||||
|
// more — once it reads below 70%, open or not.
|
||||||
|
func TestAFillingBucketOrLogIsRaisedAtThreeQuartersAndClearedBelowSevenTenths(t *testing.T) {
|
||||||
|
const mib = 1024 * 1024
|
||||||
|
log := bucketOrLog{Module: "mesh-issues", Kind: "log", Name: "changes", Stream: "LOG_mesh-issues_changes"}
|
||||||
|
bucket := bucketOrLog{Module: "mesh-issues", Kind: "bucket", Name: "issues", Stream: "KV_mesh-issues_issues"}
|
||||||
|
for _, c := range []struct {
|
||||||
|
name string
|
||||||
|
of bucketOrLog
|
||||||
|
bytes uint64
|
||||||
|
max uint64
|
||||||
|
open bool
|
||||||
|
said bool
|
||||||
|
}{
|
||||||
|
{"a log at exactly 75%", log, 768 * mib, 1024 * mib, false, true},
|
||||||
|
{"a bucket at exactly 75%", bucket, 48 * mib, 64 * mib, false, true},
|
||||||
|
{"a bucket full", bucket, 64 * mib, 64 * mib, false, true},
|
||||||
|
{"a log just under 75%", log, 768*mib - 1, 1024 * mib, false, false},
|
||||||
|
{"a bucket at 72%, not open", bucket, 64 * mib * 72 / 100, 64 * mib, false, false},
|
||||||
|
{"a bucket at 72%, open: kept", bucket, 64 * mib * 72 / 100, 64 * mib, true, true},
|
||||||
|
{"a log at exactly 70%, open: kept", log, 700 * mib, 1000 * mib, true, true},
|
||||||
|
{"a log just under 70%, open: cleared", log, 700*mib - 1, 1000 * mib, true, false},
|
||||||
|
{"a bucket at 10%, open: cleared", bucket, 64 * mib / 10, 64 * mib, true, false},
|
||||||
|
{"a bucket with no cap", bucket, 100, 0, true, false},
|
||||||
|
} {
|
||||||
|
o, said := fillObservation(c.of, filled{Bytes: c.bytes, Max: c.max}, c.open)
|
||||||
|
if said != c.said {
|
||||||
|
t.Errorf("%s: said %v, want %v", c.name, said, c.said)
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
if !said {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
if o.Key() != fillKey(c.of) || o.Kind != kindBucketOrLogFilling || o.Severity != conditions.Warning {
|
||||||
|
t.Errorf("%s: raised as %s (%s, %s)", c.name, o.Key(), o.Kind, o.Severity)
|
||||||
|
}
|
||||||
|
for _, part := range []string{c.of.Module, c.of.Kind + " " + c.of.Name, "%", " of "} {
|
||||||
|
if !strings.Contains(o.Summary, part) {
|
||||||
|
t.Errorf("%s: does not name %q: %q", c.name, part, o.Summary)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if !strings.Contains(o.Said, "bytes") {
|
||||||
|
t.Errorf("%s: the evidence does not say the fill in bytes: %q", c.name, o.Said)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
o, _ := fillObservation(log, filled{Bytes: 800 * mib, Max: 1024 * mib}, false)
|
||||||
|
if !strings.Contains(o.Summary, "800.0 MiB of 1024.0 MiB, 78%") {
|
||||||
|
t.Errorf("the fill is said as %q", o.Summary)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// **One bucket or log that cannot be read does not stop the others being judged** (review of #231): it is
|
||||||
|
// said as unread with its reason, a fill condition open for it is kept, and every other bucket and log
|
||||||
|
// is judged by its own fill.
|
||||||
|
func TestAnUnreadableBucketOrLogIsSaidAndTheRestAreStillJudged(t *testing.T) {
|
||||||
|
const mib = 1024 * 1024
|
||||||
|
broken := bucketOrLog{Module: "a", Kind: "bucket", Name: "broken", Stream: "KV_a_broken"}
|
||||||
|
brokenOpen := bucketOrLog{Module: "a", Kind: "log", Name: "was-filling", Stream: "LOG_a_was-filling"}
|
||||||
|
full := bucketOrLog{Module: "b", Kind: "log", Name: "changes", Stream: "LOG_b_changes"}
|
||||||
|
empty := bucketOrLog{Module: "c", Kind: "bucket", Name: "quiet", Stream: "KV_c_quiet"}
|
||||||
|
absent := bucketOrLog{Module: "d", Kind: "bucket", Name: "not-yet", Stream: "KV_d_not-yet"}
|
||||||
|
refusal := errors.New("the bus did not answer")
|
||||||
|
read := func(_ context.Context, s bucketOrLog) (filled, bool, error) {
|
||||||
|
switch s {
|
||||||
|
case broken, brokenOpen:
|
||||||
|
return filled{}, false, refusal
|
||||||
|
case full:
|
||||||
|
return filled{Bytes: 900 * mib, Max: 1000 * mib}, true, nil
|
||||||
|
case empty:
|
||||||
|
return filled{Bytes: 1, Max: 64 * mib}, true, nil
|
||||||
|
}
|
||||||
|
return filled{}, false, nil
|
||||||
|
}
|
||||||
|
open := func(_ context.Context, key string) (conditions.Condition, bool, error) {
|
||||||
|
if key == fillKey(brokenOpen) {
|
||||||
|
return conditions.Condition{Key: key, Severity: conditions.Warning,
|
||||||
|
Summary: "a's log was-filling holds 800.0 MiB of 1024.0 MiB, 78%"}, true, nil
|
||||||
|
}
|
||||||
|
return conditions.Condition{}, false, nil
|
||||||
|
}
|
||||||
|
got := map[string]conditions.Observation{}
|
||||||
|
for _, o := range judgeFills(context.Background(), []bucketOrLog{broken, brokenOpen, full, empty, absent}, read, open) {
|
||||||
|
got[o.Key()] = o
|
||||||
|
}
|
||||||
|
for _, s := range []bucketOrLog{broken, brokenOpen} {
|
||||||
|
o, said := got[conditions.Key(conditions.ScopeBus, s.Stream, "unread")]
|
||||||
|
if !said || o.Kind != kindBucketOrLogUnread || !o.Confirm || !strings.Contains(o.Said, refusal.Error()) {
|
||||||
|
t.Errorf("%s could not be read and was said as %+v", s.Stream, o)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if o, kept := got[fillKey(brokenOpen)]; !kept || !strings.Contains(o.Said, "kept open") {
|
||||||
|
t.Errorf("an open fill condition of a log not read this time was not kept: %+v", o)
|
||||||
|
}
|
||||||
|
if _, said := got[fillKey(broken)]; said {
|
||||||
|
t.Error("a bucket not read and not filling was said filling")
|
||||||
|
}
|
||||||
|
if o, said := got[fillKey(full)]; !said || o.Kind != kindBucketOrLogFilling {
|
||||||
|
t.Errorf("a log at 90%% beside an unreadable one was not judged: %+v", got)
|
||||||
|
}
|
||||||
|
if len(got) != 4 {
|
||||||
|
t.Errorf("said %d findings, want 4 (two unread, one kept, one filling): %v", len(got), got)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// Both kinds have plain words of their own, which hold to the plain rule (novox/hq ADR 0253).
|
||||||
|
func TestAFillingOrUnreadBucketOrLogIsSaidInPlainWords(t *testing.T) {
|
||||||
|
for _, kind := range []string{kindBucketOrLogFilling, kindBucketOrLogUnread} {
|
||||||
|
wording, has := plainWordings[kind]
|
||||||
|
if !has {
|
||||||
|
t.Errorf("%s has no plain words", kind)
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
if why, ok := conditions.PlainWords(wording(conditions.Observation{Kind: kind})); !ok {
|
||||||
|
t.Errorf("%s: %s", kind, why)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// **When neither the bucket or log nor its condition can be read, nothing is said of its fill** (second
|
||||||
|
// review of #231): the unread finding covers it, and no fill is invented for it. For one that is read and
|
||||||
|
// measured between seven tenths and three quarters, a condition that cannot be read is taken as open, so
|
||||||
|
// an unknown never clears it.
|
||||||
|
func TestNothingIsSaidOfAFillWhenNeitherItNorItsConditionCanBeRead(t *testing.T) {
|
||||||
|
const mib = 1024 * 1024
|
||||||
|
broken := bucketOrLog{Module: "a", Kind: "log", Name: "changes", Stream: "LOG_a_changes"}
|
||||||
|
between := bucketOrLog{Module: "b", Kind: "bucket", Name: "issues", Stream: "KV_b_issues"}
|
||||||
|
read := func(_ context.Context, s bucketOrLog) (filled, bool, error) {
|
||||||
|
if s == broken {
|
||||||
|
return filled{}, false, errors.New("the bus did not answer")
|
||||||
|
}
|
||||||
|
return filled{Bytes: 72 * mib, Max: 100 * mib}, true, nil
|
||||||
|
}
|
||||||
|
open := func(context.Context, string) (conditions.Condition, bool, error) {
|
||||||
|
return conditions.Condition{}, false, errors.New("the conditions bucket did not answer")
|
||||||
|
}
|
||||||
|
got := map[string]conditions.Observation{}
|
||||||
|
for _, o := range judgeFills(context.Background(), []bucketOrLog{broken, between}, read, open) {
|
||||||
|
got[o.Key()] = o
|
||||||
|
}
|
||||||
|
if o, said := got[fillKey(broken)]; said {
|
||||||
|
t.Fatalf("a fill was said for a log whose fill and condition could not be read: %+v", o)
|
||||||
|
}
|
||||||
|
if _, said := got[conditions.Key(conditions.ScopeBus, broken.Stream, "unread")]; !said {
|
||||||
|
t.Error("the log that could not be read was not said as unread")
|
||||||
|
}
|
||||||
|
if _, kept := got[fillKey(between)]; !kept {
|
||||||
|
t.Error("a bucket at 72% whose condition could not be read was cleared")
|
||||||
|
}
|
||||||
|
if len(got) != 2 {
|
||||||
|
t.Errorf("said %d findings, want 2: %v", len(got), got)
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -179,6 +179,14 @@ func moduleCheckFor(paths []string, longestMachine int, out io.Writer) error {
|
|||||||
}
|
}
|
||||||
fmt.Fprintf(out, ", keeps state %s", strings.Join(kept, ", "))
|
fmt.Fprintf(out, ", keeps state %s", strings.Join(kept, ", "))
|
||||||
}
|
}
|
||||||
|
// And the logs it keeps, with their caps (novox/hq ADR 0297).
|
||||||
|
if len(m.Logs) > 0 {
|
||||||
|
kept := make([]string, 0, len(m.Logs))
|
||||||
|
for _, l := range m.Logs {
|
||||||
|
kept = append(kept, fmt.Sprintf("%s (%d MiB)", l.Name, l.Cap()))
|
||||||
|
}
|
||||||
|
fmt.Fprintf(out, ", keeps log %s", strings.Join(kept, ", "))
|
||||||
|
}
|
||||||
if len(m.Reads) > 0 {
|
if len(m.Reads) > 0 {
|
||||||
fmt.Fprintf(out, ", reads %s", strings.Join(m.Reads, ", "))
|
fmt.Fprintf(out, ", reads %s", strings.Join(m.Reads, ", "))
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -136,6 +136,12 @@ var probeRegistry = []probe{
|
|||||||
{ID: "D-root", Asserts: "no agent can become root without a person on a machine where the router or a channel " +
|
{ID: "D-root", Asserts: "no agent can become root without a person on a machine where the router or a channel " +
|
||||||
"proving its sender runs: not by its own account, and not through a tool that runs its command as an account " +
|
"proving its sender runs: not by its own account, and not through a tool that runs its command as an account " +
|
||||||
"that can", From: "ADR 0259 §8", Kind: kindRootNotFree, Phase: 2, run: probeAgentRoot},
|
"that can", From: "ADR 0259 §8", Kind: kindRootNotFree, Phase: 2, run: probeAgentRoot},
|
||||||
|
// A module's bucket or log filling toward its cap (novox/hq ADR 0297 §6): said at three quarters, cleared
|
||||||
|
// below seven tenths, so one at its threshold is not raised and cleared on every run. One that cannot be
|
||||||
|
// read is said on its own, and every other is still judged.
|
||||||
|
{ID: probeFillID, Asserts: "every module's bucket and log holds less than three quarters of its cap; one " +
|
||||||
|
"said filling is cleared once it holds less than seven tenths", From: "ADR 0297, issue 501",
|
||||||
|
Kind: kindBucketOrLogFilling, Raises: []string{kindBucketOrLogUnread}, Phase: 1, run: probeFill},
|
||||||
{ID: "DW", Asserts: "the watchdogs of the signals table ran within three of their intervals",
|
{ID: "DW", Asserts: "the watchdogs of the signals table ran within three of their intervals",
|
||||||
From: "ADR 0227 rule 6: the watchers are watched", Kind: "watchdogs-silent", Phase: 1, run: probeWatchdogs},
|
From: "ADR 0227 rule 6: the watchers are watched", Kind: "watchdogs-silent", Phase: 1, run: probeWatchdogs},
|
||||||
// The core's health definitions (novox/hq to-be 45 §8, ADR 0236): what a core component's new build is
|
// The core's health definitions (novox/hq to-be 45 §8, ADR 0236): what a core component's new build is
|
||||||
|
|||||||
@@ -606,6 +606,19 @@ var plainWordings = map[string]func(conditions.Observation) words{
|
|||||||
Explanation: "The bus refused messages from a part of the mesh, so what they carried did not happen.",
|
Explanation: "The bus refused messages from a part of the mesh, so what they carried did not happen.",
|
||||||
Resolved: "Resolved: the bus takes the messages again"}
|
Resolved: "Resolved: the bus takes the messages again"}
|
||||||
}),
|
}),
|
||||||
|
kindBucketOrLogFilling: worded(func(o conditions.Observation) words {
|
||||||
|
return words{Headline: "A module's bucket or log on the bus is filling up",
|
||||||
|
Needs: "decide whether to raise its cap or have the module keep less.",
|
||||||
|
Explanation: "A bucket or log a module keeps on the bus holds three quarters of its cap or more. When it " +
|
||||||
|
"is full, the bus refuses what the module writes there.",
|
||||||
|
Resolved: "It has room again"}
|
||||||
|
}),
|
||||||
|
kindBucketOrLogUnread: worded(func(o conditions.Observation) words {
|
||||||
|
return words{Headline: "A module's bucket or log could not be read",
|
||||||
|
Explanation: "The controller could not read how full a bucket or log a module keeps on the bus is, so it " +
|
||||||
|
"cannot say whether it is filling up. It asks again at its next check.",
|
||||||
|
Resolved: "It can be read again"}
|
||||||
|
}),
|
||||||
"stream-wrong": worded(func(o conditions.Observation) words {
|
"stream-wrong": worded(func(o conditions.Observation) words {
|
||||||
return words{Headline: "Part of the bus's storage is wrong",
|
return words{Headline: "Part of the bus's storage is wrong",
|
||||||
Needs: "check the machine the bus runs on; the details say what is missing.",
|
Needs: "check the machine the bus runs on; the details say what is missing.",
|
||||||
|
|||||||
@@ -1300,6 +1300,15 @@ func issueMemberships(ctx context.Context, open *stores, server *link.Server, se
|
|||||||
if _, err := broker.RaiseBuckets(broker.OnConn(bus.Conn), buckets); err != nil {
|
if _, err := broker.RaiseBuckets(broker.OnConn(bus.Conn), buckets); err != nil {
|
||||||
return fmt.Errorf("the modules' state could not be asserted on the bus: %w", err)
|
return fmt.Errorf("the modules' state could not be asserted on the bus: %w", err)
|
||||||
}
|
}
|
||||||
|
// **And every declared log, for the same reason** (novox/hq ADR 0297 §2): a membership names its
|
||||||
|
// logs, and a module whose log does not exist fails its first append.
|
||||||
|
logs, err := open.inventory.DeclaredLogs(ctx)
|
||||||
|
if err != nil {
|
||||||
|
return fmt.Errorf("the modules' logs could not be read, so no log was asserted: %w", err)
|
||||||
|
}
|
||||||
|
if _, err := broker.RaiseLogs(broker.OnConn(bus.Conn), logs); err != nil {
|
||||||
|
return fmt.Errorf("the modules' logs could not be asserted on the bus: %w", err)
|
||||||
|
}
|
||||||
// Every membership is tried, and the first failure named once.
|
// Every membership is tried, and the first failure named once.
|
||||||
issued := 0
|
issued := 0
|
||||||
refused := map[string]error{}
|
refused := map[string]error{}
|
||||||
@@ -1483,13 +1492,28 @@ func raiseTheBus(ctx context.Context, inv *inventory.Inventory, address string)
|
|||||||
fmt.Printf("the bus holds state nothing declares any more, kept because it is data: %s — "+
|
fmt.Printf("the bus holds state nothing declares any more, kept because it is data: %s — "+
|
||||||
"removing it is a person's act\n", strings.Join(undeclared, ", "))
|
"removing it is a person's act\n", strings.Join(undeclared, ", "))
|
||||||
}
|
}
|
||||||
|
// Every module's log (novox/hq ADR 0297), from the catalogue, as its buckets: one that nothing
|
||||||
|
// declares any more is said and kept — a log is a module's record, and no path of the mesh removes it.
|
||||||
|
logs, err := inv.DeclaredLogs(ctx)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
unlogged, err := broker.RaiseLogs(js, logs)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
if len(unlogged) > 0 {
|
||||||
|
fmt.Printf("the bus holds logs nothing declares any more, kept because they are data: %s — "+
|
||||||
|
"removing one is a person's act\n", strings.Join(unlogged, ", "))
|
||||||
|
}
|
||||||
// And how every module hears what it consumes: asserted with the rest above, counted here.
|
// And how every module hears what it consumes: asserted with the rest above, counted here.
|
||||||
hearing, err := moduleConsumerCount(ctx, inv)
|
hearing, err := moduleConsumerCount(ctx, inv)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
fmt.Printf("the bus at %s has its streams, %d machine(s) can hear a declaration, %d module(s) "+
|
fmt.Printf("the bus at %s has its streams, %d machine(s) can hear a declaration, %d module(s) "+
|
||||||
"can hear what they consume, and %d bucket(s) of state\n", broker.BareAddress(address), len(names), hearing, len(buckets))
|
"can hear what they consume, %d bucket(s) of state and %d log(s)\n", broker.BareAddress(address), len(names), hearing,
|
||||||
|
len(buckets), len(logs))
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -448,3 +448,41 @@ func (j *JetStream) BucketNames() ([]string, error) {
|
|||||||
}
|
}
|
||||||
return out, lister.Error()
|
return out, lister.Error()
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// EnsureLog creates a module's log if it is absent and brings its configuration to match if it is
|
||||||
|
// present (novox/hq ADR 0297).
|
||||||
|
//
|
||||||
|
// **An update, never a delete and recreate**, for a bucket's reason: recreating discards what the log
|
||||||
|
// holds, and a log is a module's record. A configuration the server will not change in place (its
|
||||||
|
// storage, say, on a stream made by hand) is said as this assertion's error and the stream is left as
|
||||||
|
// it is — never removed to be made again.
|
||||||
|
func (j *JetStream) EnsureLog(l Log) error {
|
||||||
|
js, err := jetstream.New(j.conn)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
|
||||||
|
defer cancel()
|
||||||
|
if _, err := js.CreateOrUpdateStream(ctx, l.Config()); err != nil {
|
||||||
|
return fmt.Errorf("asserting log %s: %w", l.Stream(), err)
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// LogStreams is every log stream on the server: every stream whose name starts with LogStreamPrefix.
|
||||||
|
func (j *JetStream) LogStreams() ([]string, error) {
|
||||||
|
js, err := jetstream.New(j.conn)
|
||||||
|
if err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
|
||||||
|
defer cancel()
|
||||||
|
lister := js.StreamNames(ctx)
|
||||||
|
var out []string
|
||||||
|
for name := range lister.Name() {
|
||||||
|
if strings.HasPrefix(name, LogStreamPrefix) {
|
||||||
|
out = append(out, name)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return out, lister.Err()
|
||||||
|
}
|
||||||
|
|||||||
@@ -0,0 +1,195 @@
|
|||||||
|
package broker
|
||||||
|
|
||||||
|
import (
|
||||||
|
"fmt"
|
||||||
|
"sort"
|
||||||
|
"strings"
|
||||||
|
|
||||||
|
"github.com/nats-io/nats.go/jetstream"
|
||||||
|
)
|
||||||
|
|
||||||
|
// A module's log on the bus (novox/hq ADR 0297): its record of operations, each entry appended under
|
||||||
|
// a key and kept as long as the log is.
|
||||||
|
//
|
||||||
|
// A log follows a bucket in every respect (ADR 0201): the module names it locally, the mesh derives
|
||||||
|
// its stream and subjects, the controller creates it from the catalogue on every raise and from
|
||||||
|
// registration, never a module, and **no path of the mesh removes one** — a log whose declaration
|
||||||
|
// is gone is reported, as a bucket is. Unlike a bucket, a log is its owner's alone: no other module
|
||||||
|
// reads it, so there is no read of a log to grant or issue.
|
||||||
|
//
|
||||||
|
// Pure, but for the stream's configuration, which is the server's own type so that what is tested
|
||||||
|
// is exactly what is sent; jetstream.go is the part that asks a server.
|
||||||
|
|
||||||
|
// The mesh's caps on a log: what an entry's value may weigh, what a message on the log's stream may
|
||||||
|
// weigh — the value and its headers, which the server counts in a message's size, so a value of the
|
||||||
|
// full size still fits with the expected-last-sequence header an append carries — and what a log holds
|
||||||
|
// when its module says nothing and at most.
|
||||||
|
const (
|
||||||
|
LogMaxEntryBytes = 256 * 1024
|
||||||
|
LogMaxMessageBytes = LogMaxEntryBytes + 4*1024
|
||||||
|
LogDefaultMiB = 1024
|
||||||
|
LogMostMiB = 8192
|
||||||
|
)
|
||||||
|
|
||||||
|
// A Log is one module's declared log as the bus holds it.
|
||||||
|
type Log struct {
|
||||||
|
Module string
|
||||||
|
Name string
|
||||||
|
// MaxMiB is its cap in MiB; zero is LogDefaultMiB.
|
||||||
|
MaxMiB int
|
||||||
|
}
|
||||||
|
|
||||||
|
// LogStreamPrefix starts every log's stream name, which no other stream of the mesh's starts with.
|
||||||
|
const LogStreamPrefix = "LOG_"
|
||||||
|
|
||||||
|
// LogStreamName is the stream a module's log lives in: `LOG_<module>_<name>`. The module and the
|
||||||
|
// local name are each one token with no underscore, so two modules can never derive one stream.
|
||||||
|
func LogStreamName(module, name string) string { return LogStreamPrefix + module + "_" + name }
|
||||||
|
|
||||||
|
// LogSubject is the subject a log's entries are published under, without the key: an entry for key K
|
||||||
|
// is on `<LogSubject>.<K>`.
|
||||||
|
func LogSubject(module, name string) string { return "mesh.log." + module + "." + name }
|
||||||
|
|
||||||
|
// Stream is this log's stream name.
|
||||||
|
func (l Log) Stream() string { return LogStreamName(l.Module, l.Name) }
|
||||||
|
|
||||||
|
// Subject is this log's subject, without the key.
|
||||||
|
func (l Log) Subject() string { return LogSubject(l.Module, l.Name) }
|
||||||
|
|
||||||
|
// MaxBytes is this log's cap in bytes.
|
||||||
|
func (l Log) MaxBytes() int64 {
|
||||||
|
mib := l.MaxMiB
|
||||||
|
if mib <= 0 {
|
||||||
|
mib = LogDefaultMiB
|
||||||
|
}
|
||||||
|
return int64(mib) * 1024 * 1024
|
||||||
|
}
|
||||||
|
|
||||||
|
// Why is carried into the server's description of the stream, so somebody reading the server's own
|
||||||
|
// state finds whose it is and why it is kept.
|
||||||
|
func (l Log) Why() string {
|
||||||
|
return fmt.Sprintf("%s's log %q (novox/hq ADR 0297): its record of operations, one entry per operation "+
|
||||||
|
"under its key, appended by %s alone and kept as long as the log; never removed by the mesh, because "+
|
||||||
|
"it is data", l.Module, l.Name, l.Module)
|
||||||
|
}
|
||||||
|
|
||||||
|
// Config is the stream a log is, field by field as novox/hq ADR 0297 fixes it: a file stream kept by
|
||||||
|
// limits, with no maximum age and no cap per key, whose message is an entry's value and 4 KiB of
|
||||||
|
// headers at most, that refuses a new entry when full rather than
|
||||||
|
// drop an old one, and that refuses deleting an entry or purging it. Direct gets are allowed, which
|
||||||
|
// is how the runtime reads it.
|
||||||
|
func (l Log) Config() jetstream.StreamConfig {
|
||||||
|
return jetstream.StreamConfig{
|
||||||
|
Name: l.Stream(),
|
||||||
|
Description: l.Why(),
|
||||||
|
Subjects: []string{l.Subject() + ".>"},
|
||||||
|
Storage: jetstream.FileStorage,
|
||||||
|
Retention: jetstream.LimitsPolicy,
|
||||||
|
MaxAge: 0,
|
||||||
|
MaxMsgs: -1,
|
||||||
|
MaxMsgsPerSubject: -1,
|
||||||
|
MaxBytes: l.MaxBytes(),
|
||||||
|
MaxMsgSize: LogMaxMessageBytes,
|
||||||
|
Discard: jetstream.DiscardNew,
|
||||||
|
AllowDirect: true,
|
||||||
|
DenyDelete: true,
|
||||||
|
DenyPurge: true,
|
||||||
|
Replicas: 1,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// LogIssued is one log an assignment may reach, by the name its module uses for it (novox/hq ADR
|
||||||
|
// 0297): its stream, the subject its entries go under, and whether it may append. Only the owner's
|
||||||
|
// instances are issued a log, so Writes is always true today; it is said so the runtime need not
|
||||||
|
// assume it.
|
||||||
|
type LogIssued struct {
|
||||||
|
Name string `json:"name"`
|
||||||
|
Stream string `json:"stream"`
|
||||||
|
Subject string `json:"subject"`
|
||||||
|
Writes bool `json:"writes"`
|
||||||
|
}
|
||||||
|
|
||||||
|
// logsIssuedFor is every log a module's code may reach, as its membership lists them: its own.
|
||||||
|
func logsIssuedFor(d Declared) []LogIssued {
|
||||||
|
var out []LogIssued
|
||||||
|
for _, l := range d.Logs {
|
||||||
|
if !safeSubject.MatchString(d.Module) || !safeSubject.MatchString(l.Name) {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
out = append(out, LogIssued{Name: l.Name, Stream: LogStreamName(d.Module, l.Name),
|
||||||
|
Subject: LogSubject(d.Module, l.Name), Writes: true})
|
||||||
|
}
|
||||||
|
return out
|
||||||
|
}
|
||||||
|
|
||||||
|
// logGrants is what a principal publishes to reach the logs its module keeps: for each, appending
|
||||||
|
// under the log's subjects, binding to its stream, and reading it directly. Nothing more — the
|
||||||
|
// runtime reads a log by direct gets alone and makes no consumer on it — and nothing of any other
|
||||||
|
// module's log, because a log is its owner's alone. Replies come on the principal's inbox, as for a
|
||||||
|
// bucket.
|
||||||
|
func logGrants(module string, names []string) []string {
|
||||||
|
if !safeSubject.MatchString(module) {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
var out []string
|
||||||
|
for _, name := range names {
|
||||||
|
if !safeSubject.MatchString(name) {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
stream := LogStreamName(module, name)
|
||||||
|
out = append(out,
|
||||||
|
LogSubject(module, name)+".>",
|
||||||
|
"$JS.API.STREAM.INFO."+stream,
|
||||||
|
"$JS.API.DIRECT.GET."+stream,
|
||||||
|
"$JS.API.DIRECT.GET."+stream+".>")
|
||||||
|
}
|
||||||
|
return out
|
||||||
|
}
|
||||||
|
|
||||||
|
// logNames is the local names of a module's logs.
|
||||||
|
func logNames(logs []Log) []string {
|
||||||
|
out := make([]string, 0, len(logs))
|
||||||
|
for _, l := range logs {
|
||||||
|
out = append(out, l.Name)
|
||||||
|
}
|
||||||
|
return out
|
||||||
|
}
|
||||||
|
|
||||||
|
// A LogAsserter is the part of a JetStream connection log assertion needs. It has no way to remove a
|
||||||
|
// log, by design: nothing the mesh runs asks for one.
|
||||||
|
type LogAsserter interface {
|
||||||
|
// EnsureLog creates the log's stream if absent and brings its configuration to match if present,
|
||||||
|
// never deleting or recreating it.
|
||||||
|
EnsureLog(l Log) error
|
||||||
|
// LogStreams is every log stream on the server: every stream named with LogStreamPrefix.
|
||||||
|
LogStreams() ([]string, error)
|
||||||
|
}
|
||||||
|
|
||||||
|
// RaiseLogs asserts every declared log and answers the log streams on the server that nothing
|
||||||
|
// declares any more.
|
||||||
|
//
|
||||||
|
// **Those are reported, never removed** (novox/hq ADR 0297 §2, ADR 0201 §15–16): a log is a module's
|
||||||
|
// record, and a manifest edited, a module renamed or unassigned is an ordinary day's work that must
|
||||||
|
// not take a record with it. Removing one is a person's act, outside the mesh.
|
||||||
|
func RaiseLogs(a LogAsserter, logs []Log) (undeclared []string, err error) {
|
||||||
|
sorted := append([]Log(nil), logs...)
|
||||||
|
sort.Slice(sorted, func(i, j int) bool { return sorted[i].Stream() < sorted[j].Stream() })
|
||||||
|
declared := map[string]bool{}
|
||||||
|
for _, l := range sorted {
|
||||||
|
if err := a.EnsureLog(l); err != nil {
|
||||||
|
return nil, fmt.Errorf("asserting %s's log %q: %w", l.Module, l.Name, err)
|
||||||
|
}
|
||||||
|
declared[l.Stream()] = true
|
||||||
|
}
|
||||||
|
names, err := a.LogStreams()
|
||||||
|
if err != nil {
|
||||||
|
return nil, fmt.Errorf("listing the bus's logs: %w", err)
|
||||||
|
}
|
||||||
|
for _, n := range names {
|
||||||
|
if strings.HasPrefix(n, LogStreamPrefix) && !declared[n] {
|
||||||
|
undeclared = append(undeclared, n)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
sort.Strings(undeclared)
|
||||||
|
return undeclared, nil
|
||||||
|
}
|
||||||
@@ -0,0 +1,294 @@
|
|||||||
|
package broker
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"encoding/json"
|
||||||
|
"go/ast"
|
||||||
|
"go/parser"
|
||||||
|
"go/token"
|
||||||
|
"io/fs"
|
||||||
|
"path/filepath"
|
||||||
|
"slices"
|
||||||
|
"strings"
|
||||||
|
"testing"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"github.com/nats-io/nats.go/jetstream"
|
||||||
|
)
|
||||||
|
|
||||||
|
// **The stream a log is, field by field** (novox/hq ADR 0297 §2, the shared contract): a file stream
|
||||||
|
// kept by limits, no maximum age, no cap per key, the declared cap, a message of at most 260 KiB (a
|
||||||
|
// value of 256 KiB and 4 KiB of headers),
|
||||||
|
// refusing what comes next when full, read directly, refusing a delete or a purge, one replica, and
|
||||||
|
// a description that says whose it is and why.
|
||||||
|
func TestALogsStreamIsAsTheDecisionFixesIt(t *testing.T) {
|
||||||
|
got := Log{Module: "mesh-issues", Name: "changes"}.Config()
|
||||||
|
want := jetstream.StreamConfig{
|
||||||
|
Name: "LOG_mesh-issues_changes",
|
||||||
|
Description: got.Description,
|
||||||
|
Subjects: []string{"mesh.log.mesh-issues.changes.>"},
|
||||||
|
Storage: jetstream.FileStorage,
|
||||||
|
Retention: jetstream.LimitsPolicy,
|
||||||
|
MaxAge: 0,
|
||||||
|
MaxMsgs: -1,
|
||||||
|
MaxMsgsPerSubject: -1,
|
||||||
|
MaxBytes: 1024 * 1024 * 1024,
|
||||||
|
MaxMsgSize: 260 * 1024,
|
||||||
|
Discard: jetstream.DiscardNew,
|
||||||
|
AllowDirect: true,
|
||||||
|
DenyDelete: true,
|
||||||
|
DenyPurge: true,
|
||||||
|
Replicas: 1,
|
||||||
|
}
|
||||||
|
a, _ := json.Marshal(got)
|
||||||
|
b, _ := json.Marshal(want)
|
||||||
|
if string(a) != string(b) {
|
||||||
|
t.Fatalf("the log's stream is\n %s\nwant\n %s", a, b)
|
||||||
|
}
|
||||||
|
if !strings.Contains(got.Description, "mesh-issues") || !strings.Contains(got.Description, "ADR 0297") {
|
||||||
|
t.Fatalf("the description does not say whose the log is and why: %q", got.Description)
|
||||||
|
}
|
||||||
|
if c := (Log{Module: "m", Name: "n", MaxMiB: 8192}).Config(); c.MaxBytes != 8192*1024*1024 {
|
||||||
|
t.Fatalf("a cap of 8192 MiB is %d bytes", c.MaxBytes)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// **The membership names each of the owner's logs** with exactly the field names the runtime reads:
|
||||||
|
// `logs`, and in each `name`, `stream`, `subject`, `writes`. A module with no log is issued none, and
|
||||||
|
// the field is absent.
|
||||||
|
func TestAMembershipListsItsModulesLogs(t *testing.T) {
|
||||||
|
m := MembershipFor("one", Declared{Module: "mesh-issues",
|
||||||
|
Logs: []Log{{Module: "mesh-issues", Name: "changes"}}}, Placements{})
|
||||||
|
raw, err := json.Marshal(m)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
var back map[string]json.RawMessage
|
||||||
|
if err := json.Unmarshal(raw, &back); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if got := string(back["logs"]); got !=
|
||||||
|
`[{"name":"changes","stream":"LOG_mesh-issues_changes","subject":"mesh.log.mesh-issues.changes","writes":true}]` {
|
||||||
|
t.Fatalf("the membership's logs are %s", got)
|
||||||
|
}
|
||||||
|
none, _ := json.Marshal(MembershipFor("one", Declared{Module: "audit"}, Placements{}))
|
||||||
|
if strings.Contains(string(none), `"logs"`) {
|
||||||
|
t.Fatalf("a module with no log was issued logs: %s", none)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// logGrantsOf is the grants of a principal that are about logs.
|
||||||
|
func logGrantsOf(publish []string) []string {
|
||||||
|
var out []string
|
||||||
|
for _, s := range publish {
|
||||||
|
if strings.HasPrefix(s, "mesh.log.") || strings.Contains(s, ".LOG_") {
|
||||||
|
out = append(out, s)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
slices.Sort(out)
|
||||||
|
return out
|
||||||
|
}
|
||||||
|
|
||||||
|
// **The runtime is granted exactly the contract's subjects for each log it carries, and nothing else
|
||||||
|
// of any log**: appending under the log's subjects, binding to its stream, reading it directly. No
|
||||||
|
// consumer, no delete, no purge, and nothing of a log its modules do not keep.
|
||||||
|
func TestTheRuntimeIsGrantedItsModulesLogsAndNoMore(t *testing.T) {
|
||||||
|
perms, err := PermissionsFor(Principal{Kind: KindNodeTools, Node: "one", Module: RuntimeModule,
|
||||||
|
Carries: []Declared{
|
||||||
|
{Module: "mesh-issues", Logs: []Log{{Module: "mesh-issues", Name: "changes"}}},
|
||||||
|
{Module: "audit"},
|
||||||
|
}, PasswordHash: "x"})
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
want := []string{
|
||||||
|
"$JS.API.DIRECT.GET.LOG_mesh-issues_changes",
|
||||||
|
"$JS.API.DIRECT.GET.LOG_mesh-issues_changes.>",
|
||||||
|
"$JS.API.STREAM.INFO.LOG_mesh-issues_changes",
|
||||||
|
"mesh.log.mesh-issues.changes.>",
|
||||||
|
}
|
||||||
|
if got := logGrantsOf(perms.Publish); !slices.Equal(got, want) {
|
||||||
|
t.Fatalf("the runtime is granted\n %q\nwant\n %q", got, want)
|
||||||
|
}
|
||||||
|
for _, s := range perms.Subscribe {
|
||||||
|
if strings.HasPrefix(s, "mesh.log.") || strings.Contains(s, "LOG_") {
|
||||||
|
t.Fatalf("the runtime subscribes a log's subjects: %q", s)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
// A module running on its own account is granted the same for its own logs.
|
||||||
|
own, err := PermissionsFor(Principal{Kind: KindModule, Node: "one", Module: "mesh-issues",
|
||||||
|
Logs: []string{"changes"}, PasswordHash: "x"})
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if got := logGrantsOf(own.Publish); !slices.Equal(got, want) {
|
||||||
|
t.Fatalf("the module's own account is granted\n %q\nwant\n %q", got, want)
|
||||||
|
}
|
||||||
|
// And a module with no log, nothing of any.
|
||||||
|
none, err := PermissionsFor(Principal{Kind: KindModule, Node: "one", Module: "audit", PasswordHash: "x"})
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if got := logGrantsOf(none.Publish); len(got) != 0 {
|
||||||
|
t.Fatalf("a module with no log is granted %q", got)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// A log name or module that is not one plain token is issued and granted nothing rather than a
|
||||||
|
// pattern that happens to parse.
|
||||||
|
func TestALogThatNamesNoStreamGrantsNothing(t *testing.T) {
|
||||||
|
if got := logGrants("a", []string{"x.y", "x>", "*", ""}); len(got) != 0 {
|
||||||
|
t.Fatalf("granted %v for logs that name no stream", got)
|
||||||
|
}
|
||||||
|
if got := logGrants("a.b", []string{"c"}); len(got) != 0 {
|
||||||
|
t.Fatalf("granted %v for a module that is not one token", got)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
type logs struct {
|
||||||
|
ensured []string
|
||||||
|
on []string
|
||||||
|
}
|
||||||
|
|
||||||
|
func (l *logs) EnsureLog(x Log) error { l.ensured = append(l.ensured, x.Stream()); return nil }
|
||||||
|
func (l *logs) LogStreams() ([]string, error) { return l.on, nil }
|
||||||
|
|
||||||
|
// Every declared log is asserted; one on the server that nothing declares is said, not removed — and
|
||||||
|
// the asserter has no way to remove one.
|
||||||
|
func TestRaisingLogsReportsWhatNothingDeclares(t *testing.T) {
|
||||||
|
l := &logs{on: []string{"LOG_mesh-issues_changes", "LOG_gone_old", "KV_not_a_log"}}
|
||||||
|
undeclared, err := RaiseLogs(l, []Log{{Module: "mesh-issues", Name: "changes"}, {Module: "a", Name: "b"}})
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if !slices.Equal(l.ensured, []string{"LOG_a_b", "LOG_mesh-issues_changes"}) {
|
||||||
|
t.Fatalf("asserted %v", l.ensured)
|
||||||
|
}
|
||||||
|
if !slices.Equal(undeclared, []string{"LOG_gone_old"}) {
|
||||||
|
t.Fatalf("reported %v", undeclared)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// **No path of the mesh deletes or purges a log or a bucket, or recreates one** (novox/hq ADR 0297 §2,
|
||||||
|
// ADR 0201 §15–16): no code of this repository outside its tests calls the client's stream or bucket
|
||||||
|
// delete, or purges a stream, but for the controller lease's own stream, which purges its own history
|
||||||
|
// below a sequence (internal/lease). A new call is a decision, not a refactor.
|
||||||
|
func TestNoPathDeletesALogOrABucket(t *testing.T) {
|
||||||
|
removers := map[string]bool{"DeleteStream": true, "DeleteKeyValue": true, "PurgeStream": true,
|
||||||
|
"DeleteObjectStore": true}
|
||||||
|
root := filepath.Join("..", "..")
|
||||||
|
var found []string
|
||||||
|
for _, dir := range []string{"cmd", "internal"} {
|
||||||
|
err := filepath.WalkDir(filepath.Join(root, dir), func(path string, e fs.DirEntry, err error) error {
|
||||||
|
if err != nil || e.IsDir() || !strings.HasSuffix(path, ".go") || strings.HasSuffix(path, "_test.go") {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
f, err := parser.ParseFile(token.NewFileSet(), path, nil, 0)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
ast.Inspect(f, func(n ast.Node) bool {
|
||||||
|
call, ok := n.(*ast.CallExpr)
|
||||||
|
if !ok {
|
||||||
|
return true
|
||||||
|
}
|
||||||
|
sel, ok := call.Fun.(*ast.SelectorExpr)
|
||||||
|
if !ok {
|
||||||
|
return true
|
||||||
|
}
|
||||||
|
name := sel.Sel.Name
|
||||||
|
if removers[name] || (name == "Purge" && !strings.Contains(filepath.ToSlash(path), "internal/lease/")) {
|
||||||
|
found = append(found, path+": "+name)
|
||||||
|
}
|
||||||
|
return true
|
||||||
|
})
|
||||||
|
for _, lit := range stringsIn(f) {
|
||||||
|
if strings.Contains(lit, "$JS.API.STREAM.DELETE") || strings.Contains(lit, "$JS.API.STREAM.PURGE") {
|
||||||
|
// Only the writers table names it, as a subject no one but the controller may publish.
|
||||||
|
if !strings.HasSuffix(filepath.ToSlash(path), "internal/broker/writers.go") {
|
||||||
|
found = append(found, path+": "+lit)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
})
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if len(found) > 0 {
|
||||||
|
t.Fatalf("a path of the mesh removes a stream, a bucket or what one holds:\n %s", strings.Join(found, "\n "))
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// stringsIn is every string literal of a file.
|
||||||
|
func stringsIn(f *ast.File) []string {
|
||||||
|
var out []string
|
||||||
|
ast.Inspect(f, func(n ast.Node) bool {
|
||||||
|
if lit, ok := n.(*ast.BasicLit); ok && lit.Kind == token.STRING {
|
||||||
|
out = append(out, lit.Value)
|
||||||
|
}
|
||||||
|
return true
|
||||||
|
})
|
||||||
|
return out
|
||||||
|
}
|
||||||
|
|
||||||
|
// Against a real server: a log is created as its configuration says, asserting it again keeps what
|
||||||
|
// it holds, a changed cap is brought to match in place, an entry cannot be deleted from it, and one
|
||||||
|
// nothing declares any more is reported and stays.
|
||||||
|
func TestALogIsAssertedInPlaceAndNeverRemoved(t *testing.T) {
|
||||||
|
bus := aLiveBus(t)
|
||||||
|
l := Log{Module: "logtest", Name: "changes", MaxMiB: 1}
|
||||||
|
if _, err := RaiseLogs(bus, []Log{l}); err != nil {
|
||||||
|
t.Fatalf("a real server refused a module's log: %v", err)
|
||||||
|
}
|
||||||
|
js, err := jetstream.New(bus.Conn())
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
|
||||||
|
defer cancel()
|
||||||
|
ack, err := js.Publish(ctx, l.Subject()+".7", []byte(`{"op":"opened"}`))
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("an append to the log was refused: %v", err)
|
||||||
|
}
|
||||||
|
stream, err := js.Stream(ctx, l.Stream())
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
have := stream.CachedInfo().Config
|
||||||
|
want := l.Config()
|
||||||
|
if have.Storage != want.Storage || have.Retention != want.Retention || have.MaxAge != 0 ||
|
||||||
|
have.MaxMsgsPerSubject != -1 || have.MaxBytes != want.MaxBytes || have.MaxMsgSize != want.MaxMsgSize ||
|
||||||
|
have.Discard != jetstream.DiscardNew || !have.AllowDirect || !have.DenyDelete || !have.DenyPurge ||
|
||||||
|
have.Replicas != 1 || !slices.Equal(have.Subjects, want.Subjects) {
|
||||||
|
t.Fatalf("the server holds the log as %+v", have)
|
||||||
|
}
|
||||||
|
if err := stream.DeleteMsg(ctx, ack.Sequence); err == nil {
|
||||||
|
t.Fatal("an entry was deleted from a log")
|
||||||
|
}
|
||||||
|
l.MaxMiB = 2
|
||||||
|
if _, err := RaiseLogs(bus, []Log{l}); err != nil {
|
||||||
|
t.Fatalf("asserting the log again failed, so a restart would: %v", err)
|
||||||
|
}
|
||||||
|
info, err := stream.Info(ctx)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if info.Config.MaxBytes != 2*1024*1024 {
|
||||||
|
t.Fatalf("the changed cap was not brought to match: %d", info.Config.MaxBytes)
|
||||||
|
}
|
||||||
|
if info.State.Msgs < 1 {
|
||||||
|
t.Fatal("asserting the log again lost what it held")
|
||||||
|
}
|
||||||
|
undeclared, err := RaiseLogs(bus, nil)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if !slices.Contains(undeclared, l.Stream()) {
|
||||||
|
t.Fatalf("a log nothing declares was not reported: %v", undeclared)
|
||||||
|
}
|
||||||
|
if _, err := js.Stream(ctx, l.Stream()); err != nil {
|
||||||
|
t.Fatalf("a log nothing declares is gone: %v", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -51,6 +51,10 @@ type Membership struct {
|
|||||||
// and refuses, with the reason, what is not on it — the bus enforces only the union over every
|
// and refuses, with the reason, what is not on it — the bus enforces only the union over every
|
||||||
// module on the machine.
|
// module on the machine.
|
||||||
State []StateIssued `json:"state,omitempty"`
|
State []StateIssued `json:"state,omitempty"`
|
||||||
|
// Logs is every log this module's code may reach, by the name it uses for each (novox/hq ADR 0297):
|
||||||
|
// its own and no other's, since a log is its owner's alone. The runtime answers a bundle's log verbs
|
||||||
|
// from this list and refuses, with the reason, a log not on it.
|
||||||
|
Logs []LogIssued `json:"logs,omitempty"`
|
||||||
// SeatTraffic is what this module's code may submit, say, hear, take, ask, answer and read on seats
|
// SeatTraffic is what this module's code may submit, say, hear, take, ask, answer and read on seats
|
||||||
// that name their caller or their kind (novox/hq ADR 0259 §3). The runtime carrying the module
|
// that name their caller or their kind (novox/hq ADR 0259 §3). The runtime carrying the module
|
||||||
// publishes, takes and answers for it only what is listed here: the bus enforces only the union
|
// publishes, takes and answers for it only what is listed here: the bus enforces only the union
|
||||||
@@ -121,6 +125,7 @@ func MembershipFor(node string, d Declared, where Placements) Membership {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
m.State = stateIssuedFor(d, node)
|
m.State = stateIssuedFor(d, node)
|
||||||
|
m.Logs = logsIssuedFor(d)
|
||||||
t := SeatTrafficOf(d.Module, d.Holds, d.Uses, d.Watches)
|
t := SeatTrafficOf(d.Module, d.Holds, d.Uses, d.Watches)
|
||||||
for _, s := range append(append([]Seat{}, d.Uses...), d.Watches...) {
|
for _, s := range append(append([]Seat{}, d.Uses...), d.Watches...) {
|
||||||
if !s.Kinded {
|
if !s.Kinded {
|
||||||
|
|||||||
@@ -132,6 +132,8 @@ type Principal struct {
|
|||||||
// no other; KeyedReads the keys of others' state it reads one key at a time (novox/hq ADR 0260).
|
// no other; KeyedReads the keys of others' state it reads one key at a time (novox/hq ADR 0260).
|
||||||
PerMachine []string
|
PerMachine []string
|
||||||
KeyedReads []KeyedRead
|
KeyedReads []KeyedRead
|
||||||
|
// Logs is the local names of the logs this principal's module keeps (novox/hq ADR 0297).
|
||||||
|
Logs []string
|
||||||
|
|
||||||
// SnapshotsTheBus is the bus's own module, the one holding mesh-broker (novox/hq ADR 0235). Its
|
// SnapshotsTheBus is the bus's own module, the one holding mesh-broker (novox/hq ADR 0235). Its
|
||||||
// whole authority is BusSnapshotGrants: it copies the streams for the night's backup and nothing
|
// whole authority is BusSnapshotGrants: it copies the streams for the night's backup and nothing
|
||||||
@@ -746,6 +748,8 @@ func PermissionsFor(p Principal) (Permissions, error) {
|
|||||||
// watched, its own written too.
|
// watched, its own written too.
|
||||||
pub = append(pub, stateGrants(stateAccess{Module: p.Module, Node: p.Node, Keeps: p.State,
|
pub = append(pub, stateGrants(stateAccess{Module: p.Module, Node: p.Node, Keeps: p.State,
|
||||||
PerMachine: p.PerMachine, Reads: p.Reads, KeyedReads: p.KeyedReads})...)
|
PerMachine: p.PerMachine, Reads: p.Reads, KeyedReads: p.KeyedReads})...)
|
||||||
|
// And its logs (novox/hq ADR 0297): appended to and read directly, its own alone.
|
||||||
|
pub = append(pub, logGrants(p.Module, p.Logs)...)
|
||||||
|
|
||||||
// 6. Its traffic on seats that name their caller or their kind, ask proofs or keep records
|
// 6. Its traffic on seats that name their caller or their kind, ask proofs or keep records
|
||||||
// (novox/hq ADR 0259 §3).
|
// (novox/hq ADR 0259 §3).
|
||||||
@@ -833,6 +837,12 @@ func PermissionsFor(p Principal) (Permissions, error) {
|
|||||||
pub = append(pub, stateGrants(stateAccess{Module: d.Module, Node: p.Node, Keeps: stateNames(d.State),
|
pub = append(pub, stateGrants(stateAccess{Module: d.Module, Node: p.Node, Keeps: stateNames(d.State),
|
||||||
PerMachine: perMachineNames(d.State), Reads: d.Reads, KeyedReads: d.KeyedReads})...)
|
PerMachine: perMachineNames(d.State), Reads: d.Reads, KeyedReads: d.KeyedReads})...)
|
||||||
}
|
}
|
||||||
|
// **And it keeps the logs of the modules it carries** (novox/hq ADR 0297): each module's own, the
|
||||||
|
// union over them. That one module's code does not append to another's log through it is the
|
||||||
|
// runtime's to keep, from the logs each membership lists.
|
||||||
|
for _, d := range p.Carries {
|
||||||
|
pub = append(pub, logGrants(d.Module, logNames(d.Logs))...)
|
||||||
|
}
|
||||||
// **Never the traffic of a trusted holder** (novox/hq ADR 0259 §8): the machine's runtime runs as the
|
// **Never the traffic of a trusted holder** (novox/hq ADR 0259 §8): the machine's runtime runs as the
|
||||||
// operator's account, which every agent runs as, so a module saying warrants or speaking for a kind
|
// operator's account, which every agent runs as, so a module saying warrants or speaking for a kind
|
||||||
// that proves its sender is never composed into it — refused here, naming it, whatever registration
|
// that proves its sender is never composed into it — refused here, naming it, whatever registration
|
||||||
|
|||||||
@@ -35,6 +35,8 @@ type Declared struct {
|
|||||||
Invokes []string
|
Invokes []string
|
||||||
// State is the state it keeps, each a bucket its instances write (novox/hq ADR 0201).
|
// State is the state it keeps, each a bucket its instances write (novox/hq ADR 0201).
|
||||||
State []Bucket
|
State []Bucket
|
||||||
|
// Logs are the logs it keeps, each a stream its instances append to (novox/hq ADR 0297).
|
||||||
|
Logs []Log
|
||||||
// Reads are other modules' state it reads, each `<module>.<name>` (novox/hq ADR 0201).
|
// Reads are other modules' state it reads, each `<module>.<name>` (novox/hq ADR 0201).
|
||||||
Reads []string
|
Reads []string
|
||||||
// KeyedReads are keys of other modules' state it reads, each one key alone: what a seat's holder is
|
// KeyedReads are keys of other modules' state it reads, each one key alone: what a seat's holder is
|
||||||
@@ -124,6 +126,7 @@ func Users(r Records) ([]Principal, error) {
|
|||||||
Holds: d.Holds, Uses: d.Uses, Watches: d.Watches, Invokes: d.Invokes,
|
Holds: d.Holds, Uses: d.Uses, Watches: d.Watches, Invokes: d.Invokes,
|
||||||
State: stateNames(d.State), Reads: d.Reads, SnapshotsTheBus: d.SnapshotsTheBus,
|
State: stateNames(d.State), Reads: d.Reads, SnapshotsTheBus: d.SnapshotsTheBus,
|
||||||
PerMachine: perMachineNames(d.State), KeyedReads: d.KeyedReads,
|
PerMachine: perMachineNames(d.State), KeyedReads: d.KeyedReads,
|
||||||
|
Logs: logNames(d.Logs),
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
if runtimeHere {
|
if runtimeHere {
|
||||||
|
|||||||
@@ -0,0 +1,105 @@
|
|||||||
|
package catalogue
|
||||||
|
|
||||||
|
import (
|
||||||
|
"bytes"
|
||||||
|
"encoding/json"
|
||||||
|
"fmt"
|
||||||
|
)
|
||||||
|
|
||||||
|
// What a module may call its logs (novox/hq ADR 0297).
|
||||||
|
//
|
||||||
|
// A log is a module's record of operations: entries appended under a key, each kept as long as the
|
||||||
|
// log is, in the order they came. A module names each log it owns **locally** — `changes`, never a
|
||||||
|
// stream or a subject (ADR 0201 §4) — and the mesh derives the stream from the module and the local
|
||||||
|
// name, as it derives a bucket. So the rule for a log's name is a state name's: one plain token, and
|
||||||
|
// a module that keeps a log has a name that is one plain token too.
|
||||||
|
|
||||||
|
// The mesh's caps on a log, in MiB: what a module may ask, and what it gets when it asks nothing.
|
||||||
|
const (
|
||||||
|
LogLeastMiB = 1
|
||||||
|
LogMostMiB = 8192
|
||||||
|
LogDefaultMiB = 1024
|
||||||
|
)
|
||||||
|
|
||||||
|
// LogDeclaration is one log a module owns: its local name, and how large it may grow.
|
||||||
|
type LogDeclaration struct {
|
||||||
|
Name string `json:"name"`
|
||||||
|
// MaxMiB is the log's cap in MiB; zero is LogDefaultMiB. When full, the log refuses new entries
|
||||||
|
// and never drops old ones.
|
||||||
|
MaxMiB int `json:"max-mib,omitempty"`
|
||||||
|
}
|
||||||
|
|
||||||
|
// Cap is the log's cap in MiB, the default where none is said.
|
||||||
|
func (l LogDeclaration) Cap() int {
|
||||||
|
if l.MaxMiB == 0 {
|
||||||
|
return LogDefaultMiB
|
||||||
|
}
|
||||||
|
return l.MaxMiB
|
||||||
|
}
|
||||||
|
|
||||||
|
// UnmarshalJSON reads a log as its bare name, or as {name, max-mib}.
|
||||||
|
func (l *LogDeclaration) UnmarshalJSON(raw []byte) error {
|
||||||
|
trimmed := bytes.TrimSpace(raw)
|
||||||
|
if len(trimmed) > 0 && trimmed[0] == '"' {
|
||||||
|
return json.Unmarshal(trimmed, &l.Name)
|
||||||
|
}
|
||||||
|
var full struct {
|
||||||
|
Name string `json:"name"`
|
||||||
|
MaxMiB *int `json:"max-mib"`
|
||||||
|
}
|
||||||
|
dec := json.NewDecoder(bytes.NewReader(trimmed))
|
||||||
|
dec.DisallowUnknownFields()
|
||||||
|
if err := dec.Decode(&full); err != nil {
|
||||||
|
return fmt.Errorf("a log is either a name or {name, max-mib}: %w", typedUnknown(err))
|
||||||
|
}
|
||||||
|
l.Name = full.Name
|
||||||
|
l.MaxMiB = 0
|
||||||
|
if full.MaxMiB != nil {
|
||||||
|
// Said, and said as nothing: refused rather than read as the default, which it did not say.
|
||||||
|
if *full.MaxMiB == 0 {
|
||||||
|
return fmt.Errorf("log %q: max-mib is between %d and %d, not 0", full.Name, LogLeastMiB, LogMostMiB)
|
||||||
|
}
|
||||||
|
l.MaxMiB = *full.MaxMiB
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// MarshalJSON writes back the short form when there is nothing else to say.
|
||||||
|
func (l LogDeclaration) MarshalJSON() ([]byte, error) {
|
||||||
|
if l.MaxMiB == 0 {
|
||||||
|
return json.Marshal(l.Name)
|
||||||
|
}
|
||||||
|
type plain LogDeclaration
|
||||||
|
return json.Marshal(plain(l))
|
||||||
|
}
|
||||||
|
|
||||||
|
// LogProblems is what is wrong with a manifest's logs.
|
||||||
|
//
|
||||||
|
// Refused at registration, for a bucket's reason: a stream name the bus cannot hold, or a cap the
|
||||||
|
// mesh would not grant, is a module that installs, starts, and is refused on its first append.
|
||||||
|
func LogProblems(m Manifest) []string {
|
||||||
|
var problems []string
|
||||||
|
if len(m.Logs) > 0 && !stateName.MatchString(m.Module) {
|
||||||
|
problems = append(problems, fmt.Sprintf(
|
||||||
|
"%s keeps a log, and a module's name is part of its logs' names, which take one plain "+
|
||||||
|
"name — no dot (novox/hq ADR 0297)", m.Module))
|
||||||
|
}
|
||||||
|
seen := map[string]bool{}
|
||||||
|
for _, l := range m.Logs {
|
||||||
|
switch {
|
||||||
|
case !stateName.MatchString(l.Name):
|
||||||
|
problems = append(problems, fmt.Sprintf(
|
||||||
|
"%s keeps log %q: a log is named locally — lower-case letters, digits and hyphens, "+
|
||||||
|
"no dot and no underscore; the mesh derives the stream (novox/hq ADR 0297)", m.Module, l.Name))
|
||||||
|
case seen[l.Name]:
|
||||||
|
problems = append(problems, fmt.Sprintf("%s keeps log %q twice", m.Module, l.Name))
|
||||||
|
}
|
||||||
|
seen[l.Name] = true
|
||||||
|
if l.MaxMiB != 0 && (l.MaxMiB < LogLeastMiB || l.MaxMiB > LogMostMiB) {
|
||||||
|
problems = append(problems, fmt.Sprintf(
|
||||||
|
"%s caps log %q at %d MiB; a log holds between %d and %d MiB (novox/hq ADR 0297)",
|
||||||
|
m.Module, l.Name, l.MaxMiB, LogLeastMiB, LogMostMiB))
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return problems
|
||||||
|
}
|
||||||
@@ -0,0 +1,68 @@
|
|||||||
|
package catalogue
|
||||||
|
|
||||||
|
import (
|
||||||
|
"encoding/json"
|
||||||
|
"strings"
|
||||||
|
"testing"
|
||||||
|
)
|
||||||
|
|
||||||
|
// A module declares the logs it keeps (novox/hq ADR 0297 §1): a log by its bare name, or with its cap
|
||||||
|
// in MiB; the default cap where none is said.
|
||||||
|
func TestAManifestMaySayWhatLogsItKeeps(t *testing.T) {
|
||||||
|
m, err := ParseManifest([]byte(`{"module":"mesh-issues","version":"1",` +
|
||||||
|
`"logs":["changes",{"name":"moves","max-mib":2048}]}`))
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if len(m.Logs) != 2 || m.Logs[0].Name != "changes" || m.Logs[1].Name != "moves" {
|
||||||
|
t.Fatalf("logs not read: %+v", m.Logs)
|
||||||
|
}
|
||||||
|
if m.Logs[0].Cap() != LogDefaultMiB || m.Logs[0].Cap() != 1024 || m.Logs[1].Cap() != 2048 {
|
||||||
|
t.Fatalf("caps read as %d and %d", m.Logs[0].Cap(), m.Logs[1].Cap())
|
||||||
|
}
|
||||||
|
out, _ := json.Marshal(m.Logs)
|
||||||
|
if string(out) != `["changes",{"name":"moves","max-mib":2048}]` {
|
||||||
|
t.Fatalf("written back as %s", out)
|
||||||
|
}
|
||||||
|
for _, edge := range []string{`{"name":"a","max-mib":1}`, `{"name":"a","max-mib":8192}`} {
|
||||||
|
if _, err := ParseManifest([]byte(`{"module":"a","version":"1","logs":[` + edge + `]}`)); err != nil {
|
||||||
|
t.Errorf("%s is inside the caps and was refused: %v", edge, err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// The catalogue check refuses a log outside the caps, with a name that is not one plain token, with a
|
||||||
|
// field the mesh does not know, or kept by a module whose own name is not one plain token.
|
||||||
|
func TestALogIsNamedLocallyAndCapped(t *testing.T) {
|
||||||
|
for _, c := range []struct{ manifest, says string }{
|
||||||
|
{`{"module":"a","version":"1","logs":["mesh.changes"]}`, `keeps log "mesh.changes": a log is named locally`},
|
||||||
|
{`{"module":"a","version":"1","logs":["my_changes"]}`, `keeps log "my_changes"`},
|
||||||
|
{`{"module":"a","version":"1","logs":["Changes"]}`, `keeps log "Changes"`},
|
||||||
|
{`{"module":"a","version":"1","logs":["c","c"]}`, `keeps log "c" twice`},
|
||||||
|
{`{"module":"a","version":"1","logs":[{"name":"c","max-mib":8193}]}`, `between 1 and 8192 MiB`},
|
||||||
|
{`{"module":"a","version":"1","logs":[{"name":"c","max-mib":-1}]}`, `between 1 and 8192 MiB`},
|
||||||
|
{`{"module":"a","version":"1","logs":[{"name":"c","max-mib":0}]}`, `max-mib is between 1 and 8192, not 0`},
|
||||||
|
{`{"module":"a","version":"1","logs":[{"name":"c","stream":"LOG_x"}]}`, `{name, max-mib}`},
|
||||||
|
{`{"module":"a.b","version":"1","logs":["c"]}`, `no dot`},
|
||||||
|
} {
|
||||||
|
_, err := ParseManifest([]byte(c.manifest))
|
||||||
|
if err == nil {
|
||||||
|
t.Errorf("%s was accepted", c.manifest)
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
if !strings.Contains(err.Error(), c.says) {
|
||||||
|
t.Errorf("%s refused for the wrong reason: %v", c.manifest, err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// **Across the whole catalogue**: every log is named locally and capped within the mesh's caps.
|
||||||
|
func TestEveryManifestsLogsAreLocalAndCapped(t *testing.T) {
|
||||||
|
var problems []string
|
||||||
|
for _, m := range theCatalogue(t) {
|
||||||
|
problems = append(problems, LogProblems(m)...)
|
||||||
|
}
|
||||||
|
if len(problems) > 0 {
|
||||||
|
t.Fatalf("the catalogue's logs are not what ADR 0297 says:\n %s", strings.Join(problems, "\n "))
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -465,6 +465,11 @@ type Manifest struct {
|
|||||||
// (novox/hq ADR 0201). Not history — that is an event — and never a secret, sealed or not.
|
// (novox/hq ADR 0201). Not history — that is an event — and never a secret, sealed or not.
|
||||||
State []StateDeclaration `json:"state,omitempty"`
|
State []StateDeclaration `json:"state,omitempty"`
|
||||||
|
|
||||||
|
// Logs are the logs of operations this module keeps on the bus, by local name: each a stream the
|
||||||
|
// controller creates and never removes, which every instance of the module appends to and reads
|
||||||
|
// (novox/hq ADR 0297). A log is its owner's alone: no other module reads it.
|
||||||
|
Logs []LogDeclaration `json:"logs,omitempty"`
|
||||||
|
|
||||||
// Settings are the defaults this module gives its settings (novox/hq ADR 0262): each key a file,
|
// Settings are the defaults this module gives its settings (novox/hq ADR 0262): each key a file,
|
||||||
// a contribution or a served fact asks for as `${setting:<key>}`, its default, and why that
|
// a contribution or a served fact asks for as `${setting:<key>}`, its default, and why that
|
||||||
// default. Only a preference has one — a font size, a width, a number of workers — and a value
|
// default. Only a preference has one — a font size, a width, a number of workers — and a value
|
||||||
@@ -1507,7 +1512,7 @@ func ParseManifest(raw []byte) (Manifest, error) {
|
|||||||
name string
|
name string
|
||||||
n int
|
n int
|
||||||
}{{"consumes", len(m.Consumes)}, {"uses", len(m.Uses)}, {"invokes", len(m.Invokes)},
|
}{{"consumes", len(m.Consumes)}, {"uses", len(m.Uses)}, {"invokes", len(m.Invokes)},
|
||||||
{"state", len(m.State)}, {"reads", len(m.Reads)}} {
|
{"state", len(m.State)}, {"logs", len(m.Logs)}, {"reads", len(m.Reads)}} {
|
||||||
if f.n > 0 {
|
if f.n > 0 {
|
||||||
said = append(said, f.name)
|
said = append(said, f.name)
|
||||||
}
|
}
|
||||||
@@ -1618,6 +1623,8 @@ func ParseManifest(raw []byte) (Manifest, error) {
|
|||||||
problems = append(problems, EventProblems(m)...)
|
problems = append(problems, EventProblems(m)...)
|
||||||
// And what it may call its state, and whose it may read (state.go, novox/hq ADR 0201).
|
// And what it may call its state, and whose it may read (state.go, novox/hq ADR 0201).
|
||||||
problems = append(problems, StateProblems(m)...)
|
problems = append(problems, StateProblems(m)...)
|
||||||
|
// And what it may call its logs, and how large it may ask them to grow (logs.go, novox/hq ADR 0297).
|
||||||
|
problems = append(problems, LogProblems(m)...)
|
||||||
// And the defaults it gives its settings (setting_defaults.go, novox/hq ADR 0262).
|
// And the defaults it gives its settings (setting_defaults.go, novox/hq ADR 0262).
|
||||||
problems = append(problems, SettingProblems(m)...)
|
problems = append(problems, SettingProblems(m)...)
|
||||||
wellFormed := true
|
wellFormed := true
|
||||||
|
|||||||
@@ -166,6 +166,8 @@ func declaredFor(m catalogue.Manifest, seats map[string]catalogue.SeatDeclaratio
|
|||||||
// And the state it keeps and reads (novox/hq ADR 0201).
|
// And the state it keeps and reads (novox/hq ADR 0201).
|
||||||
State: bucketsOf(m),
|
State: bucketsOf(m),
|
||||||
Reads: m.Reads,
|
Reads: m.Reads,
|
||||||
|
// And the logs it keeps (novox/hq ADR 0297).
|
||||||
|
Logs: logsOf(m),
|
||||||
// And the tools its health asks (novox/hq ADR 0240): the machine's node-engine is granted them.
|
// And the tools its health asks (novox/hq ADR 0240): the machine's node-engine is granted them.
|
||||||
Checks: catalogue.HealthChecks(m),
|
Checks: catalogue.HealthChecks(m),
|
||||||
// And whether it runs as an account of its own (novox/hq ADR 0259 §8).
|
// And whether it runs as an account of its own (novox/hq ADR 0259 §8).
|
||||||
@@ -222,6 +224,29 @@ func (i *Inventory) DeclaredBuckets(ctx context.Context) ([]broker.Bucket, error
|
|||||||
return out, nil
|
return out, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// logsOf is the logs a module keeps, as the bus holds them.
|
||||||
|
func logsOf(m catalogue.Manifest) []broker.Log {
|
||||||
|
var out []broker.Log
|
||||||
|
for _, l := range m.Logs {
|
||||||
|
out = append(out, broker.Log{Module: m.Module, Name: l.Name, MaxMiB: l.Cap()})
|
||||||
|
}
|
||||||
|
return out
|
||||||
|
}
|
||||||
|
|
||||||
|
// DeclaredLogs is every log the catalogue declares, registered modules assigned or not: a log exists
|
||||||
|
// from registration, as a bucket does (novox/hq ADR 0297 §2, ADR 0201 §6).
|
||||||
|
func (i *Inventory) DeclaredLogs(ctx context.Context) ([]broker.Log, error) {
|
||||||
|
declared, err := i.Catalogue(ctx)
|
||||||
|
if err != nil {
|
||||||
|
return nil, fmt.Errorf("cannot read the catalogue: %w", err)
|
||||||
|
}
|
||||||
|
var out []broker.Log
|
||||||
|
for _, m := range declared {
|
||||||
|
out = append(out, logsOf(m)...)
|
||||||
|
}
|
||||||
|
return out, nil
|
||||||
|
}
|
||||||
|
|
||||||
func asSeat(s catalogue.SeatDeclaration, declarer string) broker.Seat {
|
func asSeat(s catalogue.SeatDeclaration, declarer string) broker.Seat {
|
||||||
seat := broker.Seat{Name: s.Name, Scope: s.Scope, Accepts: s.Accepts, Emits: s.Emits,
|
seat := broker.Seat{Name: s.Name, Scope: s.Scope, Accepts: s.Accepts, Emits: s.Emits,
|
||||||
Serves: catalogue.VerbNames(s.Serves), Kinded: s.Kinded, ByCaller: s.ByCaller, Proofs: s.Proofs,
|
Serves: catalogue.VerbNames(s.Serves), Kinded: s.Kinded, ByCaller: s.ByCaller, Proofs: s.Proofs,
|
||||||
|
|||||||
Reference in New Issue
Block a user