Compare commits

..
Author SHA1 Message Date
mesh-admin 798738cc12 Merge pull request 'A shrink is read against the item's own path: a moved directory starts its size history again (hq issue 368)' (#204) from fix/368-data-shrank-path-change into main 2026-10-10 15:56:12 +00:00
jochen babfef4f3a Keep one reading when two paths are measured at one moment, rather than failing the machine's record (hq issue 368 review)
mesh/merge-gate pass: builds build-agent, mesh-controller → ace, g14, novox, shanks; no bus step; every machine composes with the change as it did without …
mesh/repo-check pass: its merge-check.sh passed
mesh/delivery delivering: 0 machine step(s) passed
2026-10-10 17:48:49 +02:00
jochen 96b0966ad8 Read a shrink against the item's own path, so a moved directory starts its size history again (hq issue 368)
mesh/merge-gate pass: builds build-agent, mesh-controller → ace, g14, novox, shanks; no bus step; every machine composes with the change as it did without …
mesh/repo-check pass: its merge-check.sh passed
mesh/delivery superseded: a newer head of the same pull request
A reading now keeps the path it was measured at, and the peak is read only from
readings at the item's path now. Before, an agent's home that moved to its own
account was compared with the operator's home it left, and data-shrank was
raised for data that was never lost. Existing readings take their item's path
now, so a shrink that is real stays raised.
2026-10-10 17:42:46 +02:00
14 changed files with 143 additions and 599 deletions
+2 -9
View File
@@ -151,13 +151,6 @@ func busCommand(ctx context.Context, args []string) error {
if len(args) > 0 && !strings.HasPrefix(args[0], "-") {
sub, args = args[0], args[1:]
}
// The view's credential, a terminal line like a person's (bus_view.go).
switch sub {
case "view-credential":
return busViewCredential(ctx, args)
case "view-revoke":
return busViewRevoke(ctx, args)
}
set := flag.NewFlagSet("bus", flag.ContinueOnError)
snapshot := set.String("snapshot-taken", "", "where the streams' snapshot a person took is, while the mesh takes none itself")
reversible := set.Bool("reversible", false, "the new version can be undone by putting the old one back")
@@ -167,14 +160,14 @@ func busCommand(ctx context.Context, args []string) error {
if rest, err := parseAround(set, args); err != nil {
return err
} else if len(rest) > 0 {
return errors.New(busUsage)
return errors.New("bus [upgrade --why … --reversible|--irreversible [--snapshot-taken <where>]]")
}
switch sub {
case "":
return busStatus(ctx)
case "upgrade":
default:
return fmt.Errorf("bus says what a bus upgrade would do, or `bus upgrade`, `bus view-credential`, `bus view-revoke` — not %q", sub)
return fmt.Errorf("bus says what a bus upgrade would do, or `bus upgrade` — not %q", sub)
}
// Everything refused before anything is done.
if err := why.require("bus upgrade"); err != nil {
-117
View File
@@ -1,117 +0,0 @@
package main
import (
"context"
"encoding/json"
"errors"
"fmt"
"net"
"strconv"
"strings"
"github.com/novox/mesh-controller/internal/broker"
"github.com/novox/mesh-controller/internal/inventory"
)
// The view's credential: the one read-only user a page in a browser connects to the bus as, over the
// bus module's WebSocket listener (novox/hq research 036, gap G1; broker.KindView).
//
// **A terminal line, like a person's credential** (operator.go): printed once, never stored — the mesh
// keeps a hash — and revoked by forgetting the row, which the next composition of the user list makes
// real. There is one view; issuing it again rotates its password.
const busUsage = "bus [upgrade --why … --reversible|--irreversible [--snapshot-taken <where>] | view-credential | view-revoke]"
// busWebSocketPort is the port the bus module's WebSocket listener is published on, mirrored from the
// nats module's manifest (its `bus-websocket` opening), because the credential names where to connect
// and the controller does not read the module's configuration. Reached across the overlay only: the
// opening is from the mesh, and the mesh's filter admits nothing else.
const busWebSocketPort = 4223
func busViewCredential(ctx context.Context, args []string) error {
if len(args) != 0 {
return errors.New("bus view-credential takes nothing: there is one view, and this prints its credential once")
}
open, err := openStores(ctx)
if err != nil {
return err
}
defer open.Close()
inv := open.inventory
// Refused here rather than at the next composition, where it would stop the whole file.
if _, err := broker.PermissionsFor(broker.Principal{Kind: broker.KindView, PasswordHash: "x"}); err != nil {
return err
}
password, err := inv.MintBusPassword(ctx, inventory.BusUser{Username: broker.ViewUser, Kind: inventory.BusView})
if err != nil {
return err
}
where, err := broker.FromEnvironment()
if err != nil && !errors.Is(err, broker.ErrNotConfigured) {
return err
}
host := where.Address
if h, _, err := net.SplitHostPort(where.Address); err == nil {
host = h
}
websocket := ""
if host != "" {
websocket = "ws://" + net.JoinHostPort(host, strconv.Itoa(busWebSocketPort))
}
held, err := json.Marshal(struct {
WebSocket string `json:"websocket,omitempty"`
URL string `json:"url,omitempty"`
Fingerprint string `json:"fingerprint,omitempty"`
User string `json:"user"`
Password string `json:"password"`
InboxPrefix string `json:"inbox_prefix"`
Hears []string `json:"hears"`
Reads string `json:"reads"`
HowToRead string `json:"how_to_read"`
}{
WebSocket: websocket, URL: "nats://" + where.Address, Fingerprint: where.Fingerprint,
User: broker.ViewUser, Password: password,
// The client must make its inboxes under the view's own prefix: its subscribe grant is
// `_INBOX.view.>` and no wider (design 25 §4), and a client's default inbox is not under it.
InboxPrefix: "_INBOX." + broker.ViewUser,
Hears: broker.ViewHears, Reads: broker.ViewBucket,
HowToRead: "direct reads only, no watch (a consumer is refused): list with a request to $JS.API.DIRECT.GET.KV_" +
broker.ViewBucket + ` carrying {"multi_last":["$KV.` + broker.ViewBucket + `.>"]}, answered until a 204 status; ` +
"read one key with $JS.API.DIRECT.GET.KV_" + broker.ViewBucket + ".$KV." + broker.ViewBucket + ".<number> " +
"(nats.js: kvm.open(bucket, {allow_direct: true}), never create); re-read the key an event's number names",
})
if err != nil {
return err
}
fmt.Printf("issued the view, which hears %s and reads the bucket %s, and nothing else\n",
strings.Join(broker.ViewHears, ", "), broker.ViewBucket)
fmt.Println(" this is the only time the credential is printed; the mesh keeps a hash")
fmt.Println(" it works once the bus has been told, which is the next push to the machine holding mesh-broker —")
fmt.Println(" and while a new build of the bus module waits for that machine, the next `bus upgrade` a person starts,")
fmt.Println(" which is also what brings the WebSocket listener it connects through")
fmt.Println()
fmt.Println(string(held))
return nil
}
func busViewRevoke(ctx context.Context, args []string) error {
if len(args) != 0 {
return errors.New("bus view-revoke takes nothing: there is one view")
}
open, err := openStores(ctx)
if err != nil {
return err
}
defer open.Close()
if err := open.inventory.ForgetBusUser(ctx, broker.ViewUser); err != nil {
return err
}
// **Revoked at the next composition, not now** — as a person is (operator revoke): the bus's users
// are a file, and the credential stops working when the file no longer names it.
fmt.Println("the view is forgotten, and its credential stops working at the next composition — " +
"push the machine holding mesh-broker to make it so")
return nil
}
+38
View File
@@ -169,6 +169,44 @@ func TestAShrinkOfMoreThanHalfIsUrgent(t *testing.T) {
}
}
// THE FALSE ALARM (issue 368), replayed through the store D13 reads: an agent's home moved from the
// operator's own home (94.7 MB) to the agent account's fresh one (490 B), and `data-shrank` was raised
// for data that was never lost. A moved item is read against its new path only, so nothing is raised —
// and a genuine shrink at the new path, a week of history later, still is.
func TestAMovedPathIsNoShrinkAndAShrinkThereStillIs(t *testing.T) {
inv := inventory.ForTest(t)
ctx := t.Context()
shelf := shelfFor(t, houseManifest)
declared := []inventory.DeclaredData{{Module: "house", Item: "config", Class: "irreplaceable", Owned: true}}
start := time.Now().Add(-6 * time.Hour)
measure := func(at time.Time, path string, size int64) []conditions.Observation {
t.Helper()
if _, err := inv.RecordData(ctx, "home", declared, map[string]map[string]inventory.Measurement{"house": {
"config": {Path: path, Size: bytesOf(size), MeasuredAt: when(at), LastWrite: when(at),
LastBackup: when(at)}}}, "", at); err != nil {
t.Fatal(err)
}
records, err := inv.Data(ctx)
if err != nil {
t.Fatal(err)
}
peaks, err := inv.DataPeaks(ctx, at.Add(-shrinkWindow))
if err != nil {
t.Fatal(err)
}
return dataFindings(records, peaks, shelf, nil, nil, at)
}
measure(start, "/home/operator/.claude", 94_700_000)
if got := findingsByKind(measure(start.Add(10*time.Minute), "/home/agent/.claude", 490)); got[kindDataShrank].Kind != "" {
t.Fatalf("a moved path raised a shrink: %+v", got[kindDataShrank])
}
measure(start.Add(2*time.Hour), "/home/agent/.claude", 300<<20)
got := findingsByKind(measure(start.Add(4*time.Hour), "/home/agent/.claude", 1<<20))[kindDataShrank]
if got.Severity != conditions.Urgent || !strings.Contains(got.Summary, "shrank") {
t.Fatalf("a genuine shrink at the new path was not raised: %+v", got)
}
}
// Data said to be written all the time and not written; data with no backup or an old one — urgent when
// irreplaceable, a warning when valuable; and a new item given its bound before it is said.
func TestQuietDataAndMissingBackupsAreSaidByClass(t *testing.T) {
+1 -67
View File
@@ -42,64 +42,8 @@ const (
// Its authority is the union of what the modules it carries would each have had for their
// tools — and nothing of what they consume, because tools are what it runs, not reactions.
KindNodeTools Kind = "node-tools"
// KindView is the one read-only principal a view onto the bus connects as (novox/hq research 036,
// gap G1): a page in a browser, over the bus module's WebSocket listener, watching the issue tracker.
// Fixed, and derived from no declaration: what it hears is ViewHears, what it reads is ViewBucket,
// and it publishes nothing but direct reads of that one bucket (ViewReads), each answered in its own
// inbox. Composed like every other user, into the same
// file, once its credential is minted (`bus view-credential`); forgotten like every other user
// (`bus view-revoke`), at the next composition.
KindView Kind = "view"
)
// ViewUser is the view's one username: there is one view, and it is nobody's machine or module.
const ViewUser = "view"
// ViewBucket is the state the view reads: the issue tracker's issues, as the bus names the bucket
// (mesh-issues's state `issues`, novox/hq ADR 0201).
var ViewBucket = BucketName("mesh-issues", "issues")
// ViewHears are the events the view subscribes, each named: the issue tracker's own, the controller's
// walks and conditions, and the delivery owner's — what a page about issues shows beside them. Subscribe
// only, and no stream or consumer of its own: a page hears what happens while it is open, and reads the
// bucket for everything before.
var ViewHears = []string{
moduleEventSubject("mesh-issues", "opened"),
moduleEventSubject("mesh-issues", "moved"),
moduleEventSubject("mesh-issues", "noted"),
moduleEventSubject("mesh-issues", "linked"),
seatEventSubject(ControllerSeat, "plan-moved"),
seatEventSubject(ControllerSeat, "condition-raised"),
seatEventSubject(ControllerSeat, "condition-changed"),
seatEventSubject(ControllerSeat, "condition-cleared"),
moduleEventSubject("mesh-delivery", "transition"),
moduleEventSubject("mesh-delivery", "group"),
}
// ViewReads are the JetStream API requests the view makes, on ViewBucket's stream and no other: binding
// (STREAM.INFO), and direct reads — one key by its subject (`DIRECT.GET.<stream>.$KV.<bucket>.<key>`), and
// the batch form on the bare subject, which answers the newest value of every key (`multi_last`) into the
// asker's inbox. The page lists the bucket with the batch, and re-reads one key when the tracker's event
// names it (every event carries the issue's `number`).
//
// **No consumer, deliberately, and so no watch.** A KV watch is a push consumer, and a push consumer's
// deliver subject is the creator's choice, delivered by the server's own client — which the server does
// not hold to the creator's permissions. Measured on 2.11.17 (2026-10-10): the view, granted
// CONSUMER.CREATE on this stream, made a consumer delivering to `mesh.mod.mesh-issues.event.opened`, and
// a module subscribed there received the bucket's entry as the tracker's event. A grant of CONSUMER.CREATE
// is a publish to any subject in the account; the view publishes nothing, so it has none (nor
// CONSUMER.DELETE, which would let it delete a module's consumer). Nothing here is a write either: no
// `$KV.<bucket>.>`, which is what a put or a delete publishes to, and no STREAM.* that defines, purges or
// deletes.
func ViewReads() []string {
stream := "KV_" + ViewBucket
return []string{
"$JS.API.STREAM.INFO." + stream,
"$JS.API.DIRECT.GET." + stream,
"$JS.API.DIRECT.GET." + stream + ".>",
}
}
// RuntimeModule is the module that IS the node's tool runtime (novox/hq ADR 0175). Where it is
// assigned, the mesh composes one runtime principal for the machine in place of that module's own,
// and the per-module containers that served tools until then stop being the way tools reach a node.
@@ -317,8 +261,6 @@ func (p Principal) Username() string {
switch p.Kind {
case KindPerson:
return "person." + p.Module
case KindView:
return ViewUser
case KindModule, KindNodeTools:
// The runtime is named exactly as the module it stands for would have been: the mesh
// issues its credential through the same path a module's takes (`module issue`), and
@@ -591,14 +533,6 @@ func PermissionsFor(p Principal) (Permissions, error) {
// itself, its replies to the asker's own inbox.
pub = append(pub, discovering()...)
case KindView:
// Hears what it is for and reads one bucket, and nothing else (ViewHears, ViewReads): no tool,
// no event of its own, no stream, no bucket written. Its requests are answered in its own
// inbox, granted below with the person's; a reply to anything is never permitted, because
// nothing is ever asked of it.
sub = append(sub, ViewHears...)
pub = append(pub, ViewReads()...)
case KindEnrolment:
// A leaked token is useless for anything but enrolling: it cannot read a declaration, hear
// an event, or subscribe any inbox but the one its own token derives (design 25 §6).
@@ -900,7 +834,7 @@ func PermissionsFor(p Principal) (Permissions, error) {
pub = unique(pub)
}
if p.Kind == KindPerson || p.Kind == KindView {
if p.Kind == KindPerson {
// An inbox to hear answers in, and nothing else. No ack subject: a person has no durable
// consumer, because nothing is delivered to a person — they ask and are answered.
sub = append(sub, p.inbox())
-2
View File
@@ -31,8 +31,6 @@ func TestTheComposedConfigMatchesTheGolden(t *testing.T) {
// The bus's own module: the snapshot API and its inbox, nothing else (novox/hq ADR 0235).
{Kind: KindModule, Node: "one", Module: "nats", SnapshotsTheBus: true,
Serves: []string{"nats_streams"}, PasswordHash: "$2a$11$bbbbbbbbbbbbbbbbbbbbbb"},
// The view: hears the issue tracker and reads its bucket, writes nothing (research 036).
{Kind: KindView, PasswordHash: "$2a$11$vvvvvvvvvvvvvvvvvvvvvv"},
})
if err != nil {
t.Fatal(err)
-4
View File
@@ -56,10 +56,6 @@ accounts {
subscribe: { allow: ["$SRV.INFO", "$SRV.INFO.shop", "$SRV.INFO.shop.>", "$SRV.PING", "$SRV.PING.shop", "$SRV.PING.shop.>", "$SRV.STATS", "$SRV.STATS.shop", "$SRV.STATS.shop.>", "_INBOX.two.shop.>", "mesh.assignment.two.shop", "mesh.mod.shop.tool.>"] }
allow_responses: { max: 1, ttl: "1m" }
} }
{ user: "view", password: "$2a$11$vvvvvvvvvvvvvvvvvvvvvv", permissions: {
publish: { allow: ["$JS.API.DIRECT.GET.KV_mesh-issues_issues", "$JS.API.DIRECT.GET.KV_mesh-issues_issues.>", "$JS.API.STREAM.INFO.KV_mesh-issues_issues"] }
subscribe: { allow: ["_INBOX.view.>", "mesh.mod.mesh-delivery.event.group", "mesh.mod.mesh-delivery.event.transition", "mesh.mod.mesh-issues.event.linked", "mesh.mod.mesh-issues.event.moved", "mesh.mod.mesh-issues.event.noted", "mesh.mod.mesh-issues.event.opened", "mesh.seat.mesh-controller.event.condition-changed", "mesh.seat.mesh-controller.event.condition-cleared", "mesh.seat.mesh-controller.event.condition-raised", "mesh.seat.mesh-controller.event.plan-moved"] }
} }
]
}
}
-6
View File
@@ -67,9 +67,6 @@ type Records struct {
Enrolling []string
// People is each person's name against the tools they may invoke, `*` for an administrator.
People map[string][]string
// View says the mesh minted the view's credential (`bus view-credential`), so the one read-only
// view principal is composed (KindView); forgotten, it is left out, like a person.
View bool
// Interchangeable is each module whose definition says its instances are the same anywhere
// (ADR 0160), which decides whether the module's plain subject is issued to every instance.
Interchangeable map[string]bool
@@ -145,9 +142,6 @@ func Users(r Records) ([]Principal, error) {
for _, person := range sortedNames(r.People) {
out = append(out, Principal{Kind: KindPerson, Module: person, Invokes: r.People[person]})
}
if r.View {
out = append(out, Principal{Kind: KindView})
}
// Refused here rather than discovered by the server. Two users with one name is a file the
// server reads as one of them, and which one depends on the order — so a module assigned to a
-230
View File
@@ -1,230 +0,0 @@
package broker
import (
"errors"
"os"
"path/filepath"
"strings"
"testing"
"time"
"github.com/nats-io/nats-server/v2/server"
"github.com/nats-io/nats.go"
"golang.org/x/crypto/bcrypt"
)
// The view against a real server, over WebSocket (novox/hq research 036): the composed user list is
// what the server reads, the listener is the shape the bus module declares (no TLS, compression on,
// reached across the overlay only), and the view with its credential binds the issue tracker's bucket,
// lists it with one batch read, follows a tracker event to re-read the key it names — and is refused
// every write: a put, a delete, an event, and a consumer delivering onto the tracker's event subject.
//
// go test ./internal/broker/ -run TestTheView
func TestTheViewReadsTheIssuesOverWebSocketAndWritesNothing(t *testing.T) {
hash := func(password string) string {
h, err := bcrypt.GenerateFromPassword([]byte(password), bcrypt.MinCost)
if err != nil {
t.Fatal(err)
}
return string(h)
}
const node = "anchor"
tracker := Principal{Kind: KindModule, Node: node, Module: "mesh-issues", State: []string{"issues"},
Emits: []string{"opened", "moved"}, PasswordHash: hash("tracker")}
accounts, err := ComposeAccounts([]Principal{
{Kind: KindController, PasswordHash: hash("controller")},
tracker,
{Kind: KindView, PasswordHash: hash("view")},
})
if err != nil {
t.Fatal(err)
}
conf := filepath.Join(t.TempDir(), "accounts.conf")
if err := os.WriteFile(conf, []byte(accounts), 0o600); err != nil {
t.Fatal(err)
}
// The server reads the composed file as the bus does — through its own parser — and listens as the
// bus module's configuration says: a WebSocket listener without TLS and with compression, beside
// the client port. Ports chosen by the system, so this runs beside a live bus.
opts, err := server.ProcessConfigFile(conf)
if err != nil {
t.Fatalf("the server refused the composed user list: %v", err)
}
opts.Host, opts.Port = "127.0.0.1", server.RANDOM_PORT
opts.JetStream, opts.StoreDir = true, t.TempDir()
opts.NoLog, opts.NoSigs = true, true
opts.Websocket = server.WebsocketOpts{Host: "127.0.0.1", Port: server.RANDOM_PORT, NoTLS: true, Compression: true}
s, err := server.NewServer(opts)
if err != nil {
t.Fatal(err)
}
go s.Start()
if !s.ReadyForConnections(30 * time.Second) {
s.Shutdown()
t.Fatal("the server did not come up")
}
t.Cleanup(func() { s.Shutdown(); s.WaitForShutdown() })
dial := func(url, user, password string, refused chan<- string) *nats.Conn {
t.Helper()
nc, err := nats.Connect(url, nats.UserInfo(user, password), nats.CustomInboxPrefix("_INBOX."+user),
nats.Compression(true), nats.ErrorHandler(func(_ *nats.Conn, _ *nats.Subscription, err error) {
if refused != nil && errors.Is(err, nats.ErrPermissionViolation) {
refused <- err.Error()
}
}))
if err != nil {
t.Fatalf("%s could not connect to %s: %v", user, url, err)
}
t.Cleanup(nc.Close)
return nc
}
// The controller defines the bucket, as it does for every module's state; the tracker writes it.
controller := dial(s.ClientURL(), "controller", "controller", nil)
cjs, _ := controller.JetStream()
if _, err := cjs.CreateKeyValue(&nats.KeyValueConfig{Bucket: ViewBucket, History: 8}); err != nil {
t.Fatalf("the controller could not define %s: %v", ViewBucket, err)
}
trackerConn := dial(s.ClientURL(), tracker.Username(), "tracker", nil)
tjs, _ := trackerConn.JetStream()
tkv, err := tjs.KeyValue(ViewBucket)
if err != nil {
t.Fatal(err)
}
if _, err := tkv.Put("365", []byte(`{"number":365,"status":"open"}`)); err != nil {
t.Fatalf("the tracker could not write its own bucket: %v", err)
}
// The view, over WebSocket with its credential.
refused := make(chan string, 8)
view := dial(s.WebsocketURL(), ViewUser, "view", refused)
if !strings.HasPrefix(view.ConnectedUrl(), "ws://") {
t.Fatalf("the view is connected to %s, not over WebSocket", view.ConnectedUrl())
}
vjs, _ := view.JetStream(nats.MaxWait(3 * time.Second))
vkv, err := vjs.KeyValue(ViewBucket)
if err != nil {
t.Fatalf("the view could not bind %s: %v", ViewBucket, err)
}
if got, err := vkv.Get("365"); err != nil {
t.Fatalf("the view could not read a key: %v", err)
} else if !strings.Contains(string(got.Value()), `"number":365`) {
t.Fatalf("the view read %q", got.Value())
}
if _, err := tkv.Put("366", []byte(`{"number":366,"status":"open"}`)); err != nil {
t.Fatal(err)
}
// The list: one batch read, the newest value of every key, into the view's own inbox, ended by the
// server's end-of-batch status (204).
listed := map[string]string{}
inbox := view.NewRespInbox()
batch, err := view.SubscribeSync(inbox)
if err != nil {
t.Fatal(err)
}
if err := view.PublishRequest("$JS.API.DIRECT.GET.KV_"+ViewBucket, inbox,
[]byte(`{"multi_last":["$KV.`+ViewBucket+`.>"]}`)); err != nil {
t.Fatal(err)
}
for {
m, err := batch.NextMsg(5 * time.Second)
if err != nil {
t.Fatalf("the view's batch read ended without its end-of-batch (%d keys so far): %v", len(listed), err)
}
if m.Header.Get("Status") == "204" {
break
}
if status := m.Header.Get("Status"); status != "" {
t.Fatalf("the batch read answered %s %s", status, m.Header.Get("Description"))
}
listed[m.Header.Get("Nats-Subject")] = string(m.Data)
}
_ = batch.Unsubscribe()
for _, key := range []string{"365", "366"} {
if !strings.Contains(listed["$KV."+ViewBucket+"."+key], `"number":`+key) {
t.Errorf("the batch read did not list %s: %v", key, listed)
}
}
// A change followed: the tracker moves 365 and says so; the view hears the event and reads the key
// it names.
moved, err := view.SubscribeSync("mesh.mod.mesh-issues.event.moved")
if err != nil {
t.Fatal(err)
}
_ = view.Flush()
if _, err := tkv.Put("365", []byte(`{"number":365,"status":"located"}`)); err != nil {
t.Fatal(err)
}
if err := trackerConn.Publish("mesh.mod.mesh-issues.event.moved", []byte(`{"number":365,"to":"located"}`)); err != nil {
t.Fatal(err)
}
if _, err := moved.NextMsg(5 * time.Second); err != nil {
t.Fatalf("the view did not hear the tracker's event: %v", err)
}
if got, err := vkv.Get("365"); err != nil || !strings.Contains(string(got.Value()), "located") {
t.Fatalf("after the event the view read %v (%v)", got, err)
}
// **The hole a watch would open, shut** (ViewReads): a consumer delivering onto the tracker's event
// subject would have the server republish the bucket there, as the tracker. Refused, and nothing
// reaches a module listening on that subject.
listener, err := trackerConn.SubscribeSync("mesh.mod.mesh-issues.event.opened")
if err != nil {
t.Fatal(err)
}
_ = trackerConn.Flush()
if _, err := vjs.AddConsumer("KV_"+ViewBucket, &nats.ConsumerConfig{Name: "w", DeliverSubject: "mesh.mod.mesh-issues.event.opened",
AckPolicy: nats.AckNonePolicy, FilterSubject: "$KV." + ViewBucket + ".>"}); err == nil {
t.Error("the view made a consumer")
}
if _, err := vkv.WatchAll(); err == nil {
t.Error("the view made a watch, which is a consumer")
}
if m, err := listener.NextMsg(2 * time.Second); err == nil {
t.Errorf("a message reached the tracker's event subject from the view: %q", m.Data)
}
// The refusals of the consumer create land as permission violations too; drained before the writes.
drain := time.After(500 * time.Millisecond)
for draining := true; draining; {
select {
case <-refused:
case <-drain:
draining = false
}
}
// And every write is refused: the server says so, and the bucket is unchanged.
if _, err := vkv.Put("367", []byte(`{"number":367}`)); err == nil {
t.Error("the view put a key")
}
if err := vkv.Delete("365"); err == nil {
t.Error("the view deleted a key")
}
if err := view.Publish("mesh.mod.mesh-issues.event.opened", []byte(`{"number":367}`)); err != nil {
t.Fatal(err)
}
_ = view.Flush()
violations := map[string]bool{}
deadline := time.After(10 * time.Second)
for len(violations) < 3 {
select {
case v := <-refused:
for _, subject := range []string{"$KV." + ViewBucket + ".367", "$KV." + ViewBucket + ".365", "mesh.mod.mesh-issues.event.opened"} {
if strings.Contains(v, subject) {
violations[subject] = true
}
}
case <-deadline:
t.Fatalf("the server refused %d of the view's 3 writes as permission violations", len(violations))
}
}
if _, err := tkv.Get("367"); !errors.Is(err, nats.ErrKeyNotFound) {
t.Errorf("after the view's put, 367 is %v", err)
}
if _, err := tkv.Get("365"); err != nil {
t.Errorf("after the view's delete, 365 is gone: %v", err)
}
}
-144
View File
@@ -1,144 +0,0 @@
package broker
import (
"reflect"
"sort"
"testing"
)
// The view (novox/hq research 036): one read-only user, composed like every other, whose whole
// authority is a list here — so a grant that is not on the list fails a test, not a review.
// Exactly what it hears, exactly what it asks, and nothing it could write or answer. A mutation that
// adds a publish grant — `$KV.<bucket>.>`, an event, a tool — fails here.
func TestTheViewHearsAndReadsAndCanPublishNothingElse(t *testing.T) {
perms, err := PermissionsFor(Principal{Kind: KindView})
if err != nil {
t.Fatal(err)
}
wantSub := append(append([]string(nil), ViewHears...), "_INBOX.view.>")
sort.Strings(wantSub)
if !reflect.DeepEqual(perms.Subscribe, wantSub) {
t.Errorf("the view subscribes\n %v\nand should subscribe exactly\n %v", perms.Subscribe, wantSub)
}
wantPub := ViewReads()
sort.Strings(wantPub)
if !reflect.DeepEqual(perms.Publish, wantPub) {
t.Errorf("the view publishes\n %v\nand should publish exactly\n %v", perms.Publish, wantPub)
}
if len(perms.PublishDeny) != 0 {
t.Errorf("the view needs no deny, because nothing it may publish reaches the controller's own: %v", perms.PublishDeny)
}
if perms.AllowResponses {
t.Error("the view may answer, and nothing is ever asked of it")
}
// Every publish grant is binding the one bucket's stream or reading it directly. **The mutation this
// holds against**: a write grant of any shape, and a consumer of any shape — a consumer's deliver
// subject is the creator's choice, so creating one is publishing anywhere (ViewReads).
stream := "KV_" + ViewBucket
for _, p := range perms.Publish {
readOnly := p == "$JS.API.STREAM.INFO."+stream ||
p == "$JS.API.DIRECT.GET."+stream ||
p == "$JS.API.DIRECT.GET."+stream+".>"
if !readOnly {
t.Errorf("the view is granted a publish on %q, which is not a read of %s", p, ViewBucket)
}
}
for _, refused := range []string{
"$KV." + ViewBucket + ".365", // a put or a delete
"$KV.mesh-controller_conditions.x", // another bucket
"$JS.API.STREAM.CREATE." + stream, // defining the stream
"$JS.API.STREAM.PURGE." + stream, // emptying it
"$JS.API.STREAM.DELETE." + stream, // deleting it
"$JS.API.STREAM.MSG.DELETE." + stream, // deleting a message
"$JS.API.CONSUMER.CREATE.KV_mesh-controller_conditions.x", // reading another bucket
"$JS.API.CONSUMER.CREATE." + stream + ".w.$KV." + ViewBucket + ".>", // a watch: delivers anywhere
"$JS.API.CONSUMER.CREATE." + stream, // an unnamed consumer
"$JS.API.CONSUMER.DELETE." + stream + ".a_mesh-issues", // a module's consumer
"$JS.FC." + stream + ".x",
"$JS.API.DIRECT.GET.KV_mesh-controller_conditions", // another bucket, directly
"$JS.API.STREAM.INFO.EVENTS", // the events stream
"$JS.API.INFO", // the account
"mesh.mod.mesh-issues.event.opened", // claiming the tracker said something
"mesh.mod.mesh-issues.tool.open", // opening an issue
"mesh.seat.issue-tracker.tool.open", // through the seat
"mesh.seat.issue-tracker.tool.open.novox", // on one machine
"mesh.seat.mesh-controller.tool.status", // the controller's verbs
"mesh.seat.mesh-controller.event.plan-moved",
"$SRV.PING",
"_INBOX.controller.x",
} {
if MayPublish(perms, refused) {
t.Errorf("the view may publish %q", refused)
}
}
for _, refused := range []string{
"mesh.mod.mesh-issues.tool.open", // a tool asked of the tracker
"mesh.mod.telegram.event.received", // another module's events
"mesh.seat.mesh-controller.event.applied",
"mesh.control.novox.report",
"_INBOX.controller.x",
"_INBOX.person.jochen.x",
"_DELIVER.controller.EVENTS",
} {
if MaySubscribe(perms, refused) {
t.Errorf("the view may subscribe %q", refused)
}
}
for _, heard := range []string{
"mesh.mod.mesh-issues.event.opened",
"mesh.mod.mesh-issues.event.moved",
"mesh.mod.mesh-issues.event.noted",
"mesh.mod.mesh-issues.event.linked",
"mesh.seat.mesh-controller.event.plan-moved",
"mesh.seat.mesh-controller.event.condition-raised",
"mesh.seat.mesh-controller.event.condition-changed",
"mesh.seat.mesh-controller.event.condition-cleared",
"mesh.mod.mesh-delivery.event.transition",
"mesh.mod.mesh-delivery.event.group",
"_INBOX.view.abc",
} {
if !MaySubscribe(perms, heard) {
t.Errorf("the view cannot subscribe %q", heard)
}
}
}
// Composed once the mesh minted its credential, and not before: its row is the whole record of it.
func TestTheViewIsComposedOnlyOnceItsCredentialIsMinted(t *testing.T) {
without, err := Users(Records{Nodes: []string{"anchor"}})
if err != nil {
t.Fatal(err)
}
for _, p := range without {
if p.Kind == KindView {
t.Fatal("the view is composed before its credential was minted")
}
}
with, err := Users(Records{Nodes: []string{"anchor"}, View: true})
if err != nil {
t.Fatal(err)
}
views := 0
for _, p := range with {
if p.Kind == KindView {
views++
if p.Username() != ViewUser {
t.Errorf("the view is called %q, and its row is %q", p.Username(), ViewUser)
}
}
}
if views != 1 {
t.Fatalf("%d view users composed; there is one view", views)
}
// And without its hash it is named as missing, like any user — never written as a user anybody is.
_, missing := WithPasswords(with, map[string]string{})
found := false
for _, m := range missing {
found = found || m == ViewUser
}
if !found {
t.Error("a view with no password was not named as missing one")
}
}
-9
View File
@@ -108,15 +108,6 @@ func (i *Inventory) BusRecords(ctx context.Context) (broker.Records, error) {
for _, p := range people {
out.People[p.Name] = p.Invokes
}
// The view is composed once its credential is minted and until it is forgotten: its row is the
// record of it, nothing else being declared about it (broker.KindView).
kept, err := i.BusUsers(ctx)
if err != nil {
return broker.Records{}, err
}
if u, minted := kept[broker.ViewUser]; minted && u.Kind == BusView {
out.View = true
}
return out, nil
}
-3
View File
@@ -45,9 +45,6 @@ const (
// BusNodeTools is a machine's tool runtime (novox/hq ADR 0175): named like the module it
// stands for, recorded as what it is.
BusNodeTools = "node-tools"
// BusView is the one read-only view onto the bus (broker.KindView): its row is the whole record of
// it, minted by `bus view-credential` and forgotten by `bus view-revoke`.
BusView = "view"
)
// MintBusPassword makes a bus password and records its hash under a username, replacing whatever was
+18 -8
View File
@@ -94,7 +94,9 @@ type DataChange struct {
}
// readingEvery is how often a measurement is kept as a reading: the shrink is read over days, and a
// row every five minutes would say the same thing sixty times an hour.
// row every five minutes would say the same thing sixty times an hour. A reading keeps the item's path
// as it stands after the measurement — the holder's, or the last one known when it named none — and an
// item at a path with no recent reading is read at once (novox/hq issue 368).
const readingEvery = 55 * time.Minute
// readingsKept is how long readings are kept.
@@ -187,10 +189,14 @@ func (i *Inventory) RecordData(ctx context.Context, machine string, declared []D
}
if m.Size != nil && m.MeasuredAt != nil && Comparable(m.Precision) {
if _, err := tx.Exec(ctx, `
insert into data_reading (machine, module, item, at, size_bytes, last_write)
select $1, $2, $3, $4, $5, $6
where not exists (select 1 from data_reading
where machine = $1 and module = $2 and item = $3 and at > $4::timestamptz - $7::interval)`,
insert into data_reading (machine, module, item, at, size_bytes, last_write, path)
select $1, $2, $3, $4, $5, $6, d.path
from data_item d
where d.machine = $1 and d.module = $2 and d.item = $3
and not exists (select 1 from data_reading
where machine = $1 and module = $2 and item = $3 and path = d.path
and at > $4::timestamptz - $7::interval)
on conflict (machine, module, item, at) do nothing`,
machine, d.Module, d.Item, *m.MeasuredAt, *m.Size, m.LastWrite,
fmt.Sprintf("%d seconds", int(readingEvery.Seconds()))); err != nil {
return change, err
@@ -281,10 +287,14 @@ func (i *Inventory) DataOf(ctx context.Context, machine, module, item string) (D
return r, err
}
// DataPeaks is each item's largest reading since a moment, keyed by DataRecord.Key.
// DataPeaks is each item's largest reading since a moment, keyed by DataRecord.Key — read only from
// readings at the item's path now (novox/hq issue 368): a directory is never compared with another one
// that once held the same item, so a moved item starts its size history again at its new path.
func (i *Inventory) DataPeaks(ctx context.Context, since time.Time) (map[string]int64, error) {
rows, err := i.store.Pool().Query(ctx, `select machine, module, item, max(size_bytes) from data_reading
where at >= $1 group by machine, module, item`, since)
rows, err := i.store.Pool().Query(ctx, `select r.machine, r.module, r.item, max(r.size_bytes)
from data_reading r
join data_item d on d.machine = r.machine and d.module = r.module and d.item = r.item and d.path = r.path
where r.at >= $1 group by r.machine, r.module, r.item`, since)
if err != nil {
return nil, err
}
+63
View File
@@ -128,3 +128,66 @@ func TestAPartialMeasurementIsNeverAReading(t *testing.T) {
t.Fatalf("%+v, %v", r, err)
}
}
// A shrink is read against what the same directory held (novox/hq issue 368): an item whose path moved
// — an agent's home moved to its own account — starts its size history again at the new path, and a
// genuine shrink at the new path is still read against what that path held.
func TestThePeakIsReadAtTheItemsPathOnly(t *testing.T) {
inv := fresh(t)
ctx := t.Context()
now := time.Date(2026, 10, 10, 0, 0, 0, 0, time.UTC)
declared := []DeclaredData{{Module: "claude-code", Item: "agent-home", Class: "valuable", Owned: true}}
measure := func(when time.Time, path string, s int64) {
t.Helper()
if _, err := inv.RecordData(ctx, "novox", declared, map[string]map[string]Measurement{"claude-code": {
"agent-home": {Path: path, Size: size(s), MeasuredAt: at(when)}}}, "", when); err != nil {
t.Fatal(err)
}
}
measure(now, "/home/operator/.claude", 94_700_000)
// The path moves ten minutes later: the new directory is measured at once, not an hour on.
measure(now.Add(10*time.Minute), "/home/agent/.claude", 490)
peaks, err := inv.DataPeaks(ctx, now.Add(-time.Hour))
if err != nil || peaks["novox/claude-code/agent-home"] != 490 {
t.Fatalf("after the path moved the peak is %v (%v), want 490: the old directory's size is no shrink "+
"of the new one", peaks, err)
}
// The new directory grows, then genuinely loses most of it: that is read against the new path's peak.
measure(now.Add(2*time.Hour), "/home/agent/.claude", 80_000_000)
measure(now.Add(4*time.Hour), "/home/agent/.claude", 1_000)
peaks, err = inv.DataPeaks(ctx, now.Add(-time.Hour))
if err != nil || peaks["novox/claude-code/agent-home"] != 80_000_000 {
t.Fatalf("a shrink at the new path is read against %v (%v), want 80000000", peaks, err)
}
// A measurement that names no path is the item's last known path's, not a new history.
measure(now.Add(6*time.Hour), "", 2_000)
peaks, err = inv.DataPeaks(ctx, now.Add(-time.Hour))
if err != nil || peaks["novox/claude-code/agent-home"] != 80_000_000 {
t.Fatalf("a measurement with no path started a new history: peak %v (%v)", peaks, err)
}
// Moving back to a directory measured before reads it against what it held then.
measure(now.Add(8*time.Hour), "/home/operator/.claude", 94_000_000)
peaks, err = inv.DataPeaks(ctx, now.Add(-time.Hour))
if err != nil || peaks["novox/claude-code/agent-home"] != 94_700_000 {
t.Fatalf("back at the first path the peak is %v (%v), want 94700000", peaks, err)
}
}
// Two measurements at the same moment at two paths keep one reading and fail nothing: the second
// would otherwise break the reading's key and lose the machine's whole record (issue 368 review).
func TestTwoPathsAtOneMomentFailNothing(t *testing.T) {
inv := fresh(t)
ctx := t.Context()
now := time.Date(2026, 10, 10, 0, 0, 0, 0, time.UTC)
declared := []DeclaredData{{Module: "claude-code", Item: "agent-home", Class: "valuable", Owned: true}}
for _, path := range []string{"/home/operator/.claude", "/home/agent/.claude"} {
if _, err := inv.RecordData(ctx, "novox", declared, map[string]map[string]Measurement{"claude-code": {
"agent-home": {Path: path, Size: size(100), MeasuredAt: at(now)}}}, "", now); err != nil {
t.Fatalf("a measurement at %s failed: %v", path, err)
}
}
r, err := inv.DataOf(ctx, "novox", "claude-code", "agent-home")
if err != nil || r.Path != "/home/agent/.claude" {
t.Fatalf("%+v, %v", r, err)
}
}
@@ -0,0 +1,21 @@
-- A reading says the path it was measured at (novox/hq issue 368).
--
-- A shrink is read against the largest reading of the last seven days. Readings were kept by machine,
-- module and item only, so when an item's path moved — on 2026-10-10 an agent's home moved from the
-- operator's own home to the agent account's (ADR 0266) — the new, fresh directory was read against
-- the old one's size, and `data-shrank` was raised for data that was never lost. A reading now keeps
-- its path, and the peak is read only from readings at the item's path now: a moved item starts its
-- size history again, and moving back to a path reads it against what that path held.
--
-- **Every reading kept so far is taken to be at its item's path now.** Where it was really measured
-- is not on record; taking the path now keeps every item's history, so a shrink that is real stays
-- raised. An item whose path moved before this migration stays compared with its old directory until
-- those readings leave the seven-day window. Readings of an item no longer kept keep no path, and are
-- read for nothing.
--
-- Numbered 0092, past 0091, the highest on main or any open branch when this was written.
alter table data_reading add column path text;
update data_reading r set path = d.path
from data_item d
where d.machine = r.machine and d.module = r.module and d.item = r.item;