The issue tracker's history lasts seven days on the events stream and its bucket fills toward a cap with no warning (issues 501, 502). A module may now declare `logs`; the controller creates each as a stream on the bus (create-or-update, never deleted or recreated, reported when undeclared), issues it to the owner's assignments in their membership, grants the runtime exactly the append and direct-read subjects, and raises a condition when any module's bucket or log reaches 75% of its cap, cleared below 70%.
295 lines
10 KiB
Go
295 lines
10 KiB
Go
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)
|
||
}
|
||
}
|