Keep calls and hand acts on the bus, answer status at once, record durations (hq to-be 45 Phase 0)

A controller restart lost every call's outcome, `status` composed the mesh
while its caller waited (18.6s live on 2026-10-06, past the 10s window), a
repair by hand left no trace, and the core's bounds had nothing measured to
be set from.

- calls: kept in the controller's bucket mesh-controller_calls (last 1000 or
  14 days, answers bounded to 64 KiB), read by id across a restart; a
  controller starting marks a stopped one's running calls abandoned; each
  call names its caller from the inbox its answer goes to.
- status: the serving controller composes it at start, after news from a
  machine, a build or an acting verb, and every minute; the verb answers the
  last composition at once with when and how long it took. Composing resolves
  each machine once instead of twice.
- hand-act log in mesh-controller_hand-acts: push (required through the seat),
  plans stop/close, broker consumer-reset and the new hand-act record take
  --why/--cause/--condition; `hand-acts` lists them and repeated causes;
  status counts the week's.
- durations (migration 0066): apply (send to first report), heartbeat gap,
  plan tier and build, recorded as heard; `durations` summarises them.
- the controller's seat row takes this binary's definition of its own verbs,
  so the console no longer judges calls against an older build's schema.
- the controller is granted its two buckets' subjects.
This commit is contained in:
jochen
2026-10-06 02:59:36 +02:00
parent 146c48fd96
commit e74c32ed50
35 changed files with 2313 additions and 75 deletions
+102
View File
@@ -0,0 +1,102 @@
package broker
import (
"context"
"fmt"
"time"
"github.com/nats-io/nats.go/jetstream"
)
// The controller's own key-value buckets (novox/hq to-be 45 §1, §6, §7).
//
// **What the controller must remember across its own restart, it keeps on the bus.** A call's
// outcome lived in the memory of the process that served it (novox/hq issue 265), so a controller
// replaced while a push ran answered "no such call" for the one thing its caller had been told to
// ask about. The bus already outlives the controller and is the shape ADR 0201 gives a module's
// current state: one value per key, written by one owner, read by anybody granted it. These are the
// controller's, written by it alone — the writers table of to-be 45 §1 — and asserted on every start
// like the streams, so a bus raised from nothing has them before the first call is served.
// CallsBucket keeps every call of the mesh's own verbs and what came of it; HandActsBucket every act
// a person did by hand, with why.
var (
CallsBucket = BucketName(ControllerSeat, "calls")
HandActsBucket = BucketName(ControllerSeat, "hand-acts")
)
// The bounds to-be 45 §6 sets for calls: the last thousand, or fourteen days, whichever is fewer.
// A call is two keys — its record, and its answer apart so a listing does not read every answer —
// so the stream holds twice as many messages as it keeps calls.
const (
KeptCallsDurably = 1000
CallsKeptFor = 14 * 24 * time.Hour
// CallAnswerBytes is the most of one answer kept: a whole declaration is far smaller, and an
// answer larger is cut and says so.
CallAnswerBytes = 64 << 10
// HandActsKeptFor is as long as a condition's history (to-be 45 §2): an act by hand is read
// back beside what it addressed.
HandActsKeptFor = 90 * 24 * time.Hour
)
// IsControllerBucket says a bucket is the controller's own, not a module's state nothing declares.
func IsControllerBucket(bucket string) bool {
return bucket == CallsBucket || bucket == HandActsBucket
}
// ControllerBucketsAsserter is what raising the controller's buckets needs of a connection.
type ControllerBucketsAsserter interface {
EnsureControllerBuckets() error
}
// EnsureControllerBuckets creates the controller's buckets if absent and brings their options to
// match. An update, never a delete: what they hold is the record of what the mesh was asked.
func (j *JetStream) EnsureControllerBuckets() 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.CreateOrUpdateKeyValue(ctx, jetstream.KeyValueConfig{
Bucket: CallsBucket,
Description: "the calls of the mesh's own verbs and what came of each (novox/hq to-be 45 §6, issue " +
"265): written by the controller alone, read through `calls`; the last thousand, or fourteen days",
History: 1,
TTL: CallsKeptFor,
MaxValueSize: CallAnswerBytes + 4<<10,
MaxBytes: 2 * KeptCallsDurably * (CallAnswerBytes + 4<<10),
Storage: jetstream.FileStorage,
}); err != nil {
return fmt.Errorf("asserting bucket %s: %w", CallsBucket, err)
}
// **The count, on the stream under the bucket.** A bucket has an age and a size and no count;
// the stream it is made of does, and with one value per key the oldest message is the oldest
// call. Asserted after the bucket, every time, because asserting the bucket writes the stream's
// configuration whole and puts the count back to none.
stream, err := js.Stream(ctx, "KV_"+CallsBucket)
if err != nil {
return fmt.Errorf("reading the stream under %s: %w", CallsBucket, err)
}
cfg := stream.CachedInfo().Config
if cfg.MaxMsgs != 2*KeptCallsDurably {
cfg.MaxMsgs = 2 * KeptCallsDurably
cfg.Discard = jetstream.DiscardOld
if _, err := js.UpdateStream(ctx, cfg); err != nil {
return fmt.Errorf("bounding %s to the last %d calls: %w", CallsBucket, KeptCallsDurably, err)
}
}
if _, err := js.CreateOrUpdateKeyValue(ctx, jetstream.KeyValueConfig{
Bucket: HandActsBucket,
Description: "every act a person did by hand, with why (novox/hq to-be 45 §7): written by the " +
"controller's repairing verbs and `hand-act record`, read through `hand-acts`",
History: 1,
TTL: HandActsKeptFor,
MaxValueSize: 16 << 10,
MaxBytes: 64 << 20,
Storage: jetstream.FileStorage,
}); err != nil {
return fmt.Errorf("asserting bucket %s: %w", HandActsBucket, err)
}
return nil
}
+5
View File
@@ -288,6 +288,11 @@ func PermissionsFor(p Principal) (Permissions, error) {
// enrolments and can reach nothing else.
pub = append(pub, "_INBOX."+enrolmentPrefix+".>")
// **And its own buckets** (novox/hq to-be 45 §1): the calls it served and the acts done by
// hand, which it alone writes. A put is a publish to the bucket's subject, which `$JS.API.>`
// does not cover; each bucket named, not `$KV.>`, which would let it write any module's state.
pub = append(pub, "$KV."+CallsBucket+".>", "$KV."+HandActsBucket+".>")
case KindPerson:
// Tools, and nothing else. Every subject a person may publish is a tool call; a person
// who could publish an event would be able to claim a module said something.
+3 -2
View File
@@ -160,8 +160,9 @@ func RaiseBuckets(a BucketAsserter, buckets []Bucket) (undeclared []string, err
return nil, fmt.Errorf("listing the bus's state: %w", err)
}
for _, n := range names {
// A seat's cancelled set is the mesh's own (novox/hq ADR 0219), not a module's state.
if !declared[n] && !IsCancelledSet(n) {
// A seat's cancelled set is the mesh's own (novox/hq ADR 0219), not a module's state; so are
// the controller's own buckets (novox/hq to-be 45 §1).
if !declared[n] && !IsCancelledSet(n) && !IsControllerBucket(n) {
undeclared = append(undeclared, n)
}
}
+1 -1
View File
@@ -24,7 +24,7 @@ accounts {
jetstream: enabled
users = [
{ user: "controller", password: "$2a$11$cccccccccccccccccccccc", permissions: {
publish: { allow: ["$JS.ACK.CONTROL.controller.>", "$JS.ACK.EVENTS.controller.>", "$JS.API.>", "_INBOX.enrol.>", "mesh.assignment.>", "mesh.control.>", "mesh.mod.*.tool.>", "mesh.node.>", "mesh.seat.mesh-build-machine.accept.>", "mesh.seat.mesh-build-machine.tool.>", "mesh.seat.mesh-controller.event.applied", "mesh.seat.mesh-controller.event.built-before", "mesh.seat.mesh-controller.event.refused", "mesh.seat.node-build-agent.accept.>", "mesh.seat.node-build-agent.tool.>"] }
publish: { allow: ["$JS.ACK.CONTROL.controller.>", "$JS.ACK.EVENTS.controller.>", "$JS.API.>", "$KV.mesh-controller_calls.>", "$KV.mesh-controller_hand-acts.>", "_INBOX.enrol.>", "mesh.assignment.>", "mesh.control.>", "mesh.mod.*.tool.>", "mesh.node.>", "mesh.seat.mesh-build-machine.accept.>", "mesh.seat.mesh-build-machine.tool.>", "mesh.seat.mesh-controller.event.applied", "mesh.seat.mesh-controller.event.built-before", "mesh.seat.mesh-controller.event.refused", "mesh.seat.node-build-agent.accept.>", "mesh.seat.node-build-agent.tool.>"] }
subscribe: { allow: ["$JS.API.>", "$SRV.INFO", "$SRV.INFO.mesh-controller", "$SRV.INFO.mesh-controller.>", "$SRV.PING", "$SRV.PING.mesh-controller", "$SRV.PING.mesh-controller.>", "$SRV.STATS", "$SRV.STATS.mesh-controller", "$SRV.STATS.mesh-controller.>", "_DELIVER.controller", "_DELIVER.controller.>", "_INBOX.controller.>", "mesh.control.>", "mesh.mod.*.event.provisioner.failing", "mesh.mod.*.event.provisioner.recovered", "mesh.mod.gitea.event.pull.merged", "mesh.mod.mesh-catalog.event.catching-up", "mesh.mod.mesh-catalog.event.upgraded", "mesh.seat.mesh-build-machine.event.built", "mesh.seat.mesh-controller.tool.>", "mesh.seat.node-build-agent.event.built"] }
allow_responses: { max: 1, ttl: "1m" }
} }
+27 -2
View File
@@ -110,6 +110,8 @@ var ControllerVerbs = []Verb{
"paths": "with repository: the files the merge would change, comma-separated, from the repository's root",
"modules": "with repository: or the modules it would change, comma-separated",
"limit": "how many plans to list (default 10); only when listing",
"why": "with stop or close: why it is ended by hand — required, and recorded in the hand-act log (novox/hq to-be 45 §7)",
"cause": "with stop or close: the cause in a word, or a condition's kind (optional)",
}, nil)},
{Name: "plan", Description: "What one machine would run, and why: the declaration the mesh would send it — " +
"or, with files, the files it would be given.",
@@ -137,11 +139,14 @@ var ControllerVerbs = []Verb{
{Name: "push", Description: "Send one machine everything it should be. With no machine named it is a push of the " +
"WHOLE mesh — every machine that is behind — and the answer says so first; behind says that outright. " +
"Answers at once that it is running, with a call id: `calls` with that id says what it sent " +
"(a push can reload the bus, which then refuses any answer still to come).",
"(a push can reload the bus, which then refuses any answer still to come). A push by hand is a repair, " +
"and says why: recorded in the hand-act log (novox/hq to-be 45 §7).",
Input: schema(map[string]string{
"node": "the machine's name; without it, every machine that is behind",
"behind": "\"true\": every machine that is behind, the whole mesh — the same as naming none, said outright; not with node",
}, nil, "behind")},
"why": "why this is pushed by hand: recorded in the hand-act log",
"cause": "the cause in a word, or a condition's kind — the word a second push for the same reason uses (optional)",
}, []string{"why"}, "behind")},
{Name: "rotate", Description: "Replace a credential. A pair credential, by provision (and a consuming machine and module, " +
"else every holder): both ends are re-sent together. Or a module's own secret, by machine, module and " +
"name: made anew and the machine sent, so the module starts again on it — only for a secret its " +
@@ -206,6 +211,26 @@ var ControllerVerbs = []Verb{
Input: schema(map[string]string{"node": "one machine; every holder when absent"}, nil)},
{Name: "resume", Description: "The build seat's holder on one machine — or every holder — takes builds again.",
Input: schema(map[string]string{"node": "one machine; every holder when absent"}, nil)},
// Acts done by hand, and what the bounds are set from (novox/hq to-be 45 §7, Phase 0).
{Name: "hand-act", Description: "Record an act done by hand outside the mesh — a container restarted, a file " +
"edited, a service started on a machine — with why and its cause, in the hand-act log beside the pushes and " +
"plans ended by hand (novox/hq to-be 45 §7). A cause recorded twice in a fortnight is a healer wanted.",
Input: schema(map[string]string{
"what": "what was done, in a line",
"why": "why it had to be done by hand",
"cause": "the cause in a word, or a condition's kind — the word a second act for the same reason uses",
"condition": "the key of the condition it addressed, if any (optional)",
}, []string{"what", "why", "cause"})},
{Name: "hand-acts", Description: "What was done by hand lately — pushes, plans ended, consumers re-made, acts " +
"recorded — who, why and the cause of each, and which causes repeat: each repeat is a healer the mesh lacks.",
Input: schema(map[string]string{"days": "how many days back (default 14)"}, nil)},
{Name: "durations", Description: "How long things take, as the controller measured them: a send to its machine's " +
"report (apply), a machine's silence between words (heartbeat-gap), a plan's tier, a build — per machine, " +
"repository or module, with median, p90 and max. What the core's bounds are set from (novox/hq to-be 45 Phase 0).",
Input: schema(map[string]string{
"kind": "one kind: apply, heartbeat-gap, plan-tier or build; every kind when absent",
"days": "how many days back (default 14)",
}, nil)},
{Name: "build", Description: "Have the build machine build a repository. Answers at once with the build's id: " +
"`builds` with that id follows it line by line, and the module is registered when the outcome comes.",
Input: schema(map[string]string{
+146
View File
@@ -0,0 +1,146 @@
package inventory
import (
"context"
"errors"
"fmt"
"time"
"github.com/jackc/pgx/v5"
)
// The durations the core's bounds are set from (novox/hq to-be 45 Phase 0): see migration 0066.
// The kinds of duration recorded.
const (
DurationApply = "apply"
DurationHeartbeatGap = "heartbeat-gap"
DurationPlanTier = "plan-tier"
DurationBuild = "build"
)
// DurationKinds are every kind, in the order `durations` shows them.
var DurationKinds = []string{DurationApply, DurationHeartbeatGap, DurationPlanTier, DurationBuild}
// DurationsKeptFor is how long a duration is kept: long enough to set a bound from, and to correct it
// in Phase 1's first live week.
const DurationsKeptFor = 30 * 24 * time.Hour
// Duration is one measurement.
type Duration struct {
Kind string `json:"kind"`
Subject string `json:"subject"`
Node string `json:"node,omitempty"`
Ref string `json:"ref"`
Started time.Time `json:"started"`
Took time.Duration `json:"took"`
Detail string `json:"detail,omitempty"`
}
// RecordDuration keeps one measurement, once: the same thing measured again is not a second row.
func (i *Inventory) RecordDuration(ctx context.Context, d Duration) error {
if d.Took < 0 {
return nil // a clock that went back measures nothing
}
_, err := i.store.Pool().Exec(ctx,
`insert into duration (kind, subject, node, ref, started, took_ms, detail)
values ($1, $2, $3, $4, $5, $6, $7) on conflict do nothing`,
d.Kind, d.Subject, d.Node, d.Ref, d.Started, d.Took.Milliseconds(), d.Detail)
return err
}
// RecordApplyDuration measures a machine's report of the declaration it was last sent: from the send
// to the first report of it. A report of anything else, or of a send already measured, measures
// nothing — a machine reconciling reports the same declaration every few minutes.
func (i *Inventory) RecordApplyDuration(ctx context.Context, node, declared, outcome string) error {
if declared == "" {
return nil
}
_, err := i.store.Pool().Exec(ctx,
`insert into duration (kind, subject, node, ref, started, took_ms, detail)
select $1, n.name, n.name, n.sent || '@' || to_char(n.sent_at at time zone 'UTC', 'YYYY-MM-DD"T"HH24:MI:SS.US'),
n.sent_at, (extract(epoch from (now() - n.sent_at)) * 1000)::bigint, $4
from node n
where n.name = $2 and n.sent = $3 and n.sent_at is not null
on conflict do nothing`,
DurationApply, node, declared, outcome)
return err
}
// RecordHeartbeatGap measures the silence before a machine's word: from the last word heard before
// it, which the caller read before recording this one.
func (i *Inventory) RecordHeartbeatGap(ctx context.Context, node string, before time.Time) error {
if before.IsZero() {
return nil
}
now := time.Now()
return i.RecordDuration(ctx, Duration{Kind: DurationHeartbeatGap, Subject: node, Node: node,
Ref: before.UTC().Format(time.RFC3339Nano), Started: before, Took: now.Sub(before)})
}
// Durations is every measurement of a kind since a moment, oldest first; every kind when kind is empty.
func (i *Inventory) Durations(ctx context.Context, kind string, since time.Time) ([]Duration, error) {
rows, err := i.store.Pool().Query(ctx,
`select kind, subject, node, ref, started, took_ms, detail from duration
where ($1 = '' or kind = $1) and recorded >= $2 order by recorded`, kind, since)
if err != nil {
return nil, err
}
defer rows.Close()
var out []Duration
for rows.Next() {
var d Duration
var ms int64
if err := rows.Scan(&d.Kind, &d.Subject, &d.Node, &d.Ref, &d.Started, &ms, &d.Detail); err != nil {
return nil, err
}
d.Took = time.Duration(ms) * time.Millisecond
out = append(out, d)
}
return out, rows.Err()
}
// ForgetOldDurations removes what is older than DurationsKeptFor, and says how many.
func (i *Inventory) ForgetOldDurations(ctx context.Context) (int64, error) {
tag, err := i.store.Pool().Exec(ctx, `delete from duration where recorded < $1`,
time.Now().Add(-DurationsKeptFor))
if err != nil {
return 0, err
}
return tag.RowsAffected(), nil
}
// planTierLeft is the measurement a plan's save makes when it leaves a tier: the tier moved on, or
// the plan ended. Read in the save's own transaction, so two saves cannot both measure one tier.
func planTierLeft(ctx context.Context, tx pgx.Tx, p Plan, now time.Time) (entered time.Time, err error) {
var oldTier int
var oldState string
var since time.Time
err = tx.QueryRow(ctx,
`select tier, state, coalesce(tier_entered, created) from release_plan where id = $1 for update`,
p.ID).Scan(&oldTier, &oldState, &since)
if errors.Is(err, pgx.ErrNoRows) {
return now, nil // a new plan enters its first tier now
}
if err != nil {
return time.Time{}, err
}
wasOpen := oldState == PlanBuilding || oldState == PlanRolling
if !wasOpen || (oldTier == p.Tier && p.Open()) {
return since, nil // still in the tier, or already ended
}
var modules []string
if oldTier >= 0 && oldTier < len(p.Tiers) {
modules = p.Tiers[oldTier]
}
detail := fmt.Sprintf("plan %s, tier %d of %d (%v), left %s", p.ID, oldTier, len(p.Tiers), modules, p.State)
if p.Open() {
detail = fmt.Sprintf("plan %s, tier %d of %d (%v), moved on to tier %d", p.ID, oldTier, len(p.Tiers), modules, p.Tier)
}
_, err = tx.Exec(ctx,
`insert into duration (kind, subject, node, ref, started, took_ms, detail)
values ($1, $2, '', $3, $4, $5, $6) on conflict do nothing`,
DurationPlanTier, p.Repository, fmt.Sprintf("%s/tier-%d", p.ID, oldTier), since,
now.Sub(since).Milliseconds(), detail)
return now, err
}
+97
View File
@@ -0,0 +1,97 @@
package inventory
import (
"testing"
"time"
)
// **The durations the bounds are set from are recorded once each** (novox/hq to-be 45 Phase 0): a
// send measured at its first report and not again at every reconcile that repeats it; a send of
// something else measures nothing; a new send is a new measurement.
func TestAnApplyIsMeasuredOncePerSend(t *testing.T) {
inv := ForTest(t)
ctx := t.Context()
node, err := inv.AddNode(ctx, "anchor")
if err != nil {
t.Fatal(err)
}
if err := inv.RecordSent(ctx, node.ID, "d1", nil); err != nil {
t.Fatal(err)
}
for range 3 { // the report, then two reconciles saying the same
if err := inv.RecordApplyDuration(ctx, "anchor", "d1", OutcomeApplied); err != nil {
t.Fatal(err)
}
}
if err := inv.RecordApplyDuration(ctx, "anchor", "d0", OutcomeApplied); err != nil {
t.Fatal(err)
}
ds, err := inv.Durations(ctx, DurationApply, time.Now().Add(-time.Hour))
if err != nil || len(ds) != 1 || ds[0].Subject != "anchor" || ds[0].Detail != OutcomeApplied || ds[0].Took < 0 {
t.Fatalf("%v %+v", err, ds)
}
time.Sleep(10 * time.Millisecond)
if err := inv.RecordSent(ctx, node.ID, "d2", nil); err != nil {
t.Fatal(err)
}
if err := inv.RecordApplyDuration(ctx, "anchor", "d2", OutcomeFailed); err != nil {
t.Fatal(err)
}
if ds, _ = inv.Durations(ctx, "", time.Now().Add(-time.Hour)); len(ds) != 2 {
t.Fatalf("a second send was not measured: %+v", ds)
}
}
// A machine's silence is the time since its last word, once per word.
func TestASilenceIsMeasuredFromTheLastWord(t *testing.T) {
inv := ForTest(t)
ctx := t.Context()
before := time.Now().Add(-90 * time.Second)
for range 2 {
if err := inv.RecordHeartbeatGap(ctx, "anchor", before); err != nil {
t.Fatal(err)
}
}
if err := inv.RecordHeartbeatGap(ctx, "anchor", time.Time{}); err != nil {
t.Fatal(err)
}
ds, err := inv.Durations(ctx, DurationHeartbeatGap, time.Now().Add(-time.Hour))
if err != nil || len(ds) != 1 || ds[0].Took < 90*time.Second || ds[0].Took > 2*time.Minute {
t.Fatalf("%v %+v", err, ds)
}
}
// A plan's tier is measured when the plan leaves it — moving on, or ending — and only then.
func TestAPlansTierIsMeasuredWhenItIsLeft(t *testing.T) {
inv := ForTest(t)
ctx := t.Context()
p := Plan{ID: "plan-1", Repository: "novox/mesh-tools", Commit: "abc", Created: time.Now().UTC(),
State: PlanBuilding, Tiers: [][]string{{"mesh-tools"}, {"builder"}},
Modules: map[string]*PlanModule{"mesh-tools": {}, "builder": {}}}
if err := inv.SavePlan(ctx, p); err != nil {
t.Fatal(err)
}
p.State = PlanRolling // the same tier, saved again
if err := inv.SavePlan(ctx, p); err != nil {
t.Fatal(err)
}
if ds, _ := inv.Durations(ctx, DurationPlanTier, time.Now().Add(-time.Hour)); len(ds) != 0 {
t.Fatalf("a tier not left was measured: %+v", ds)
}
p.Tier = 1
if err := inv.SavePlan(ctx, p); err != nil {
t.Fatal(err)
}
p.State = PlanDone
if err := inv.SavePlan(ctx, p); err != nil {
t.Fatal(err)
}
if err := inv.SavePlan(ctx, p); err != nil { // saved again once ended: nothing more to measure
t.Fatal(err)
}
ds, err := inv.Durations(ctx, DurationPlanTier, time.Now().Add(-time.Hour))
if err != nil || len(ds) != 2 || ds[0].Subject != "novox/mesh-tools" || ds[0].Ref != "plan-1/tier-0" ||
ds[1].Ref != "plan-1/tier-1" {
t.Fatalf("%v %+v", err, ds)
}
}
@@ -0,0 +1,31 @@
-- The durations the core's bounds are set from (novox/hq to-be 45 Phase 0).
--
-- Every watchdog of the core's signals table has a bound — how long a machine may take to report
-- after a send, how long it may be silent, how long a plan's tier or a build may take — and a bound
-- guessed is a condition that cries wolf or one that never fires. So the controller records each
-- duration as it is observed, and Phase 1 sets the bounds from what was recorded:
--
-- apply a declaration sent → the machine's first report of that declaration, per machine
-- heartbeat-gap one word from a machine → the next, per machine
-- plan-tier a plan entering a tier → leaving it, per repository
-- build a build asked → its outcome heard, per module
--
-- One row per thing measured: `ref` names it (the send, the earlier word, the plan's tier, the
-- build), so a report repeated by a reconcile is not a second measurement. Kept a month; read
-- through `durations`.
create table duration (
kind text not null,
subject text not null,
node text not null default '',
ref text not null,
started timestamptz not null,
took_ms bigint not null,
detail text not null default '',
recorded timestamptz not null default now(),
primary key (kind, subject, ref)
);
create index duration_by_kind on duration (kind, recorded);
-- When a plan entered the tier it is at, so leaving it measures the tier.
alter table release_plan add column tier_entered timestamptz;
+22 -6
View File
@@ -80,13 +80,29 @@ func (i *Inventory) SavePlan(ctx context.Context, p Plan) error {
if err != nil {
return err
}
_, err = i.store.Pool().Exec(ctx,
`insert into release_plan (id, repository, commit_hash, created, updated, state, tier, tiers, modules, note, branch)
values ($1, $2, $3, $4, now(), $5, $6, $7, $8, $9, $10)
// **And how long the tier it left took** (novox/hq to-be 45 Phase 0): measured here, where the
// plan moves, in the same transaction as the move, so no save can move a tier unmeasured or
// measure one twice.
tx, err := i.store.Pool().Begin(ctx)
if err != nil {
return err
}
defer func() { _ = tx.Rollback(ctx) }()
entered, err := planTierLeft(ctx, tx, p, time.Now())
if err != nil {
return err
}
_, err = tx.Exec(ctx,
`insert into release_plan (id, repository, commit_hash, created, updated, state, tier, tiers, modules, note, branch, tier_entered)
values ($1, $2, $3, $4, now(), $5, $6, $7, $8, $9, $10, $11)
on conflict (id) do update set updated = now(), state = excluded.state, tier = excluded.tier,
tiers = excluded.tiers, modules = excluded.modules, note = excluded.note, branch = excluded.branch`,
p.ID, p.Repository, p.Commit, p.Created, p.State, p.Tier, tiers, modules, p.Note, p.Branch)
return err
tiers = excluded.tiers, modules = excluded.modules, note = excluded.note, branch = excluded.branch,
tier_entered = excluded.tier_entered`,
p.ID, p.Repository, p.Commit, p.Created, p.State, p.Tier, tiers, modules, p.Note, p.Branch, entered)
if err != nil {
return err
}
return tx.Commit(ctx)
}
// OpenPlans is every plan still being worked, oldest first.
+32 -4
View File
@@ -105,14 +105,27 @@ func (i *Inventory) widenProtocol(ctx context.Context, s catalogue.Seat) error {
changed := false
row.Accepts, changed = union(row.Accepts, s.Accepts, changed)
row.Emits, changed = union(row.Emits, s.Emits, changed)
have := map[string]bool{}
for _, v := range row.Serves {
have[v.Name] = true
have := map[string]int{}
for n, v := range row.Serves {
have[v.Name] = n
}
for _, v := range s.Serves {
if !have[v.Name] {
at, kept := have[v.Name]
if !kept {
row.Serves = append(row.Serves, v)
changed = true
continue
}
// **The controller's own verbs are described by the binary that runs them** (novox/hq
// issue 244, to-be 45 §7). Its seat's verbs are not an operator's to reshape: each is a
// command line this binary composes from the arguments its own table declares, and the
// console judges a call against the row. A row kept from an older build described `push`
// without the `behind` and `why` the binary takes, so the console refused an argument the verb
// needs. So a verb this binary defines takes this binary's definition; a verb only the row has
// — a newer build's, during a roll-out (ADR 0185) — is left as it is.
if s.Name == catalogue.ControllerSeatName && !sameVerb(row.Serves[at], v) {
row.Serves[at] = v
changed = true
}
}
if !changed {
@@ -282,3 +295,18 @@ func (i *Inventory) Holdings(ctx context.Context) ([]catalogue.Held, error) {
}
return out, rows.Err()
}
// sameVerb is whether two definitions of a verb say the same, read as the row stores them.
func sameVerb(a, b catalogue.Verb) bool {
ja, errA := json.Marshal(a)
jb, errB := json.Marshal(b)
if errA != nil || errB != nil {
return false
}
var ra, rb any
_ = json.Unmarshal(ja, &ra)
_ = json.Unmarshal(jb, &rb)
ca, _ := json.Marshal(ra)
cb, _ := json.Marshal(rb)
return string(ca) == string(cb)
}
+43
View File
@@ -113,3 +113,46 @@ func TestRenameSeatKeepsTheFormerNameAsAnAlias(t *testing.T) {
t.Fatal("renaming a seat to its own name was accepted")
}
}
// **The controller's verbs in the row are the binary's** (novox/hq issue 244, to-be 45 §7): a row
// seeded by an older build describes `push` without the arguments this one takes, and the console
// judges a call against the row — so re-seeding brings the controller's verbs to this binary's
// definition, and keeps a verb only the row has, which a newer build added (ADR 0185).
func TestTheControllersVerbsInTheRowAreTheBinarys(t *testing.T) {
inv := ForTest(t)
ctx := t.Context()
if _, err := inv.SeedSeats(ctx, catalogue.DefaultSeats()); err != nil {
t.Fatal(err)
}
old := `[{"name":"push","description":"an older push","input":{"type":"object","properties":{"node":{"type":"string"}}}},
{"name":"newer","description":"a verb of a newer build"}]`
if _, err := inv.store.Pool().Exec(ctx, `update seat set serves = $1 where name = $2`,
[]byte(old), catalogue.ControllerSeatName); err != nil {
t.Fatal(err)
}
if _, err := inv.SeedSeats(ctx, catalogue.DefaultSeats()); err != nil {
t.Fatal(err)
}
seats, err := inv.Seats(ctx)
if err != nil {
t.Fatal(err)
}
verbs := map[string]catalogue.Verb{}
for _, s := range seats {
if s.Name == catalogue.ControllerSeatName {
for _, v := range s.Serves {
verbs[v.Name] = v
}
}
}
props, _ := verbs["push"].Input["properties"].(map[string]any)
if _, takesWhy := props["why"]; !takesWhy || verbs["push"].Description == "an older push" {
t.Fatalf("push in the row is still the older build's: %+v", verbs["push"])
}
if _, kept := verbs["newer"]; !kept {
t.Fatal("a verb only the row has was dropped")
}
if _, added := verbs["durations"]; !added {
t.Fatal("a verb this binary adds was not added")
}
}
+205 -17
View File
@@ -7,7 +7,9 @@ import (
"fmt"
"log"
"regexp"
"sort"
"strconv"
"strings"
"sync"
"time"
@@ -33,8 +35,9 @@ import (
// either. A variable so a test need not wait.
var AnswerWithin = 10 * time.Second
// KeptCalls is how many calls are kept, newest first; keptAnswer the largest answer kept of a call
// whose caller was sent it. One its caller never had is kept whole.
// KeptCalls is how many calls this process keeps in memory, newest first — to match a refusal to its
// call, and to answer at once; the bus keeps the last thousand (Durably). keptAnswer is the largest
// answer kept of a call whose caller was sent it. One its caller never had is kept whole.
const (
KeptCalls = 100
keptAnswer = 64 << 10
@@ -47,6 +50,10 @@ const (
// CallFinishedAfter is a call that finished after its caller was told it was still running: its
// answer is here and nowhere else.
CallFinishedAfter = "finished after its caller was answered"
// CallAbandoned is a call whose controller stopped before it finished (novox/hq to-be 45 §6): a
// controller starting finds it running under another and says so, rather than leaving it running
// for ever in the record.
CallAbandoned = "abandoned: the controller running it stopped before it finished"
)
// Call is one call of a role's tool, as `calls` shows it.
@@ -65,10 +72,27 @@ type Call struct {
// Refused is the bus refusing the answer this holder sent: the caller got nothing, and this is
// the only place that says what it would have.
Refused string `json:"answer refused by the bus,omitempty"`
// Caller is the bus principal that asked, read from the inbox its answer went to — every principal
// is granted only its own (novox/hq to-be 45 §7: a hand act says who).
Caller string `json:"caller,omitempty"`
// Holder is the controller process that served it, so one starting can tell its own running
// calls from those a stopped one left.
Holder string `json:"holder,omitempty"`
reply string // the subject the answer went to, which the bus names when it refuses it
}
// CallKeeper keeps calls where the process serving them does not: the controller's bucket on the
// bus (novox/hq to-be 45 §6). A call is kept whole on every change — begun, finished, refused — so a
// controller replaced at any moment leaves the last word on each.
type CallKeeper interface {
Keep(ctx context.Context, c Call) error
// Kept is one call with its whole answer.
Kept(ctx context.Context, id string) (Call, bool, error)
// Recent is the kept calls, newest first, without their answers.
Recent(ctx context.Context) ([]Call, error)
}
// CallLog keeps the latest calls a holder served.
type CallLog struct {
mu sync.Mutex
@@ -78,6 +102,15 @@ type CallLog struct {
// Follow is the tool that reads this log back, named in a running answer — set by a holder that
// serves one (the controller's `calls`); without it, the answer points at the holder's journal.
Follow string
// keeper keeps every call beyond this process, when Durably was given one; writes go through
// one goroutine, in order, so a call's last state is the one kept.
keeper CallKeeper
holder string
writes chan Call
logger *log.Logger
lost int // writes the keeper could not take, said once each
keepErr error
}
// Calls is this process's log: one holder process serves its seats on one connection.
@@ -85,16 +118,130 @@ var Calls = NewCallLog()
func NewCallLog() *CallLog { return &CallLog{now: time.Now} }
// keepTries is how many times one call's state is offered to the keeper before it is said lost: the
// bus reloading its user list refuses for a moment, and that is exactly when a push runs.
const keepTries = 5
// Durably keeps every call from now on with keeper as well as in memory, under this process's name,
// and marks running the calls a controller before this one left running: it stopped, so they cannot
// finish (novox/hq to-be 45 §6). Said, naming each.
func (l *CallLog) Durably(ctx context.Context, keeper CallKeeper, holder string, logger *log.Logger) error {
l.mu.Lock()
l.keeper, l.holder, l.logger = keeper, holder, logger
if l.writes == nil {
l.writes = make(chan Call, 256)
go l.keepWrites()
}
l.mu.Unlock()
kept, err := keeper.Recent(ctx)
if err != nil {
return fmt.Errorf("reading the calls kept on the bus: %w", err)
}
for _, c := range kept {
if c.State != CallRunning || c.Holder == holder {
continue
}
whole, found, err := keeper.Kept(ctx, c.ID)
if err != nil || !found {
whole = c
}
whole.State = CallAbandoned
if err := keeper.Keep(ctx, whole); err != nil {
return fmt.Errorf("marking %s abandoned: %w", c.ID, err)
}
if logger != nil {
logger.Printf("%s (%s.%s, asked %s by %s) was running under %s, which stopped: marked abandoned — "+
"it may have done part of what it was asked, and nothing will finish it", c.ID, c.Seat, c.Verb,
c.Started.Format(time.RFC3339), orSomebody(c.Caller), orSomebody(c.Holder))
}
}
return nil
}
func orSomebody(s string) string {
if s == "" {
return "an unnamed caller"
}
return s
}
// keep queues one call's state for the keeper. Never blocks a call: a queue that is full is a keeper
// that is not taking writes, and that is said rather than waited on.
func (l *CallLog) keep(c Call) {
if l.writes == nil {
return
}
select {
case l.writes <- c:
default:
l.lost++
if l.logger != nil {
l.logger.Printf("%s (%s.%s) is kept in memory only: the bus is not taking calls' records (%d not kept)",
c.ID, c.Seat, c.Verb, l.lost)
}
}
}
func (l *CallLog) keepWrites() {
for c := range l.writes {
var err error
for try := 0; try < keepTries; try++ {
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
err = l.keeper.Keep(ctx, c)
cancel()
if err == nil {
break
}
time.Sleep(time.Duration(try+1) * time.Second)
}
if err != nil && l.logger != nil {
l.logger.Printf("%s (%s.%s, %s) could not be kept on the bus after %d tries: %v — `calls` "+
"answers it from memory until this controller stops", c.ID, c.Seat, c.Verb, c.State, keepTries, err)
}
}
}
// callerOf is the bus principal an answer goes to: every principal's inbox is `_INBOX.<its user>.`
// followed by the client's own random token, and a user may itself hold dots.
func callerOf(reply string) string {
rest, ok := strings.CutPrefix(reply, "_INBOX.")
if !ok {
return ""
}
tokens := strings.Split(rest, ".")
for i, t := range tokens {
if i > 0 && isNUID(t) {
return strings.Join(tokens[:i], ".")
}
}
return ""
}
// isNUID is the client library's random inbox token: twenty-two letters and digits.
func isNUID(t string) bool {
if len(t) != 22 {
return false
}
for _, r := range t {
if !(r >= '0' && r <= '9' || r >= 'a' && r <= 'z' || r >= 'A' && r <= 'Z') {
return false
}
}
return true
}
func (l *CallLog) begin(seat, verb string, args json.RawMessage, reply string) *Call {
l.mu.Lock()
defer l.mu.Unlock()
l.next++
c := &Call{ID: "call-" + strconv.FormatInt(l.now().UnixNano(), 10) + "-" + strconv.FormatUint(l.next, 10),
Seat: seat, Verb: verb, Args: args, Started: l.now(), State: CallRunning, reply: reply}
Seat: seat, Verb: verb, Args: args, Started: l.now(), State: CallRunning, reply: reply,
Caller: callerOf(reply), Holder: l.holder}
l.calls = append(l.calls, c)
if len(l.calls) > KeptCalls {
l.calls = l.calls[len(l.calls)-KeptCalls:]
}
l.keep(*c)
return c
}
@@ -117,29 +264,67 @@ func (l *CallLog) finish(c *Call, answer []byte, failed, answeredAlready bool) {
} else {
c.State = CallAnswered
}
l.keep(*c)
}
// Recent is the kept calls, newest first, as copies.
func (l *CallLog) Recent() []Call {
// Recent is the kept calls, newest first, as copies: this process's from memory, and, when they are
// kept durably, every other the bus holds — a call a controller before this one served included.
// Memory wins for a call in both, being the newer word on it. A bus that cannot be read is said in
// the error beside what memory holds, never answered as no calls.
func (l *CallLog) Recent() ([]Call, error) {
l.mu.Lock()
defer l.mu.Unlock()
out := make([]Call, 0, len(l.calls))
seen := map[string]bool{}
for i := len(l.calls) - 1; i >= 0; i-- {
out = append(out, *l.calls[i])
seen[l.calls[i].ID] = true
}
return out
}
// Get is one kept call.
func (l *CallLog) Get(id string) (Call, bool) {
l.mu.Lock()
defer l.mu.Unlock()
for _, c := range l.calls {
if c.ID == id {
return *c, true
keeper := l.keeper
l.mu.Unlock()
if keeper == nil {
return out, nil
}
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
kept, err := keeper.Recent(ctx)
if err != nil {
return out, fmt.Errorf("the calls kept on the bus could not be read, so only this controller's own "+
"are listed: %w", err)
}
for _, c := range kept {
if !seen[c.ID] {
out = append(out, c)
}
}
return Call{}, false
sort.SliceStable(out, func(i, j int) bool { return out[i].Started.After(out[j].Started) })
return out, nil
}
// Get is one kept call: from memory, or from the bus when this process did not serve it.
func (l *CallLog) Get(id string) (Call, bool, error) {
l.mu.Lock()
for _, c := range l.calls {
if c.ID == id {
found := *c
l.mu.Unlock()
return found, true, nil
}
}
keeper := l.keeper
l.mu.Unlock()
if keeper == nil {
return Call{}, false, nil
}
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
return keeper.Kept(ctx, id)
}
// IsDurable says whether calls outlive this process.
func (l *CallLog) IsDurable() bool {
l.mu.Lock()
defer l.mu.Unlock()
return l.keeper != nil
}
// kept is what a call's arguments are kept as: each argument by name, a value only when it is short
@@ -195,6 +380,7 @@ func (l *CallLog) Refusal(err error, logger *log.Logger) bool {
at := l.now()
hit.Refused = fmt.Sprintf("%s: %s", at.Format(time.RFC3339), err)
id, verb, took := hit.ID, hit.Seat+"."+hit.Verb, at.Sub(hit.Started).Round(time.Second)
l.keep(*hit)
l.mu.Unlock()
if logger != nil {
// Said in the mesh's words, beside the library's own line: which call, and where its answer is.
@@ -277,6 +463,8 @@ func (l *CallLog) serveCall(seat, verb string, args json.RawMessage, reply strin
args = json.RawMessage(`{}`)
}
c := l.begin(seat, verb, kept(args), reply)
// Who asked travels with the call, so an act it does by hand says so (novox/hq to-be 45 §7).
ctx = context.WithValue(ctx, callerKey{}, c.Caller)
acknowledged := make(chan struct{})
var once sync.Once
ctx = context.WithValue(ctx, answerNowKey{}, func() { once.Do(func() { close(acknowledged) }) })
+115
View File
@@ -0,0 +1,115 @@
package link
import (
"context"
"encoding/json"
"errors"
"fmt"
"sort"
"github.com/nats-io/nats.go"
"github.com/nats-io/nats.go/jetstream"
"github.com/novox/mesh-controller/internal/broker"
)
// Calls kept on the bus (novox/hq to-be 45 §6).
//
// **Two keys a call**: `<id>` holds the record — verb, arguments as kept, caller, state, times — and
// `<id>.answer` the answer, bounded. A listing watches the records alone, so reading the last thousand
// calls reads the last thousand small records and not a thousand answers; one call asked by id reads
// both. A call's id is one token, so the record's key never holds a dot and `*` matches records only.
// BusCalls keeps calls in the controller's calls bucket.
type BusCalls struct {
kv jetstream.KeyValue
}
// CallsOnTheBus opens the calls bucket the controller asserts at its start.
func CallsOnTheBus(ctx context.Context, conn *nats.Conn) (*BusCalls, error) {
api, err := jetstream.New(conn)
if err != nil {
return nil, err
}
kv, err := api.KeyValue(ctx, broker.CallsBucket)
if err != nil {
return nil, fmt.Errorf("the calls bucket %s is not on the bus — the controller asserts it at its "+
"start, so one older than this has not: %w", broker.CallsBucket, err)
}
return &BusCalls{kv: kv}, nil
}
const answerKey = ".answer"
// Keep writes a call's record, and its answer when it has one.
func (b *BusCalls) Keep(ctx context.Context, c Call) error {
answer := c.Answer
c.Answer = nil
if len(answer) > 0 {
if len(answer) > broker.CallAnswerBytes {
answer, _ = json.Marshal(map[string]any{"cut": fmt.Sprintf("an answer of %d bytes; the first %d "+
"are kept", len(answer), broker.CallAnswerBytes), "start": string(answer[:broker.CallAnswerBytes])})
}
// The answer before the record, so a record saying a call finished never points at an answer
// not yet written.
if _, err := b.kv.Put(ctx, c.ID+answerKey, answer); err != nil {
return err
}
}
record, err := json.Marshal(c)
if err != nil {
return err
}
_, err = b.kv.Put(ctx, c.ID, record)
return err
}
// Kept is one call, with its answer.
func (b *BusCalls) Kept(ctx context.Context, id string) (Call, bool, error) {
entry, err := b.kv.Get(ctx, id)
if errors.Is(err, jetstream.ErrKeyNotFound) || errors.Is(err, jetstream.ErrInvalidKey) {
return Call{}, false, nil
}
if err != nil {
return Call{}, false, err
}
var c Call
if err := json.Unmarshal(entry.Value(), &c); err != nil {
return Call{}, false, fmt.Errorf("the record of %s on the bus is not a call: %w", id, err)
}
answer, err := b.kv.Get(ctx, id+answerKey)
switch {
case err == nil:
c.Answer = json.RawMessage(answer.Value())
case !errors.Is(err, jetstream.ErrKeyNotFound):
return Call{}, false, err
}
return c, true, nil
}
// Recent is every call the bucket holds, newest first, without answers: read once through a watch
// of the records, which hands over the current value of each and then says it has.
func (b *BusCalls) Recent(ctx context.Context) ([]Call, error) {
w, err := b.kv.Watch(ctx, "*", jetstream.IgnoreDeletes())
if err != nil {
return nil, err
}
defer func() { _ = w.Stop() }()
var out []Call
for {
select {
case <-ctx.Done():
return nil, fmt.Errorf("reading the calls bucket: %w", ctx.Err())
case entry := <-w.Updates():
if entry == nil {
// Every current value handed over.
sort.SliceStable(out, func(i, j int) bool { return out[i].Started.After(out[j].Started) })
return out, nil
}
var c Call
if json.Unmarshal(entry.Value(), &c) == nil && c.ID != "" {
out = append(out, c)
}
}
}
}
+152
View File
@@ -0,0 +1,152 @@
package link
import (
"bytes"
"context"
"encoding/json"
"log"
"strings"
"testing"
"time"
"github.com/nats-io/nats.go/jetstream"
"github.com/novox/mesh-controller/internal/broker"
)
// The calls bucket against a real server (novox/hq to-be 45 §6): whether a call's outcome is
// readable by id from a controller that did not serve it is a claim about what the bus keeps.
func callsBucket(t *testing.T) *BusCalls {
t.Helper()
js := aBus(t)
api, err := jetstream.New(js.Conn())
if err != nil {
t.Fatal(err)
}
_ = api.DeleteKeyValue(t.Context(), broker.CallsBucket)
if err := js.EnsureControllerBuckets(); err != nil {
t.Fatal(err)
}
keeper, err := CallsOnTheBus(t.Context(), js.Conn())
if err != nil {
t.Fatal(err)
}
return keeper
}
// waitKept waits until the keeper holds a call in a state, the writes being queued.
func waitKept(t *testing.T, keeper CallKeeper, id, state string) Call {
t.Helper()
deadline := time.Now().Add(5 * time.Second)
for {
c, found, err := keeper.Kept(context.Background(), id)
if err == nil && found && c.State == state {
return c
}
if time.Now().After(deadline) {
t.Fatalf("%s never kept as %q: %+v %v %v", id, state, c, found, err)
}
time.Sleep(20 * time.Millisecond)
}
}
// **A controller restart keeps every call's outcome** (to-be 45 Phase 0): a call one controller
// answered is read whole by id from the next, and a call it left running is said abandoned rather
// than running for ever.
func TestNatsACallsOutcomeOutlivesItsController(t *testing.T) {
keeper := callsBucket(t)
first := NewCallLog()
if err := first.Durably(t.Context(), keeper, "controller@one", nil); err != nil {
t.Fatal(err)
}
a := newAnswers(t)
first.serveCall("mesh-controller", "push", json.RawMessage(`{"node":"anchor","values":"secret"}`),
"_INBOX.node-tools.g14.Pe3sGzAtv8jUBYQ6sKSvKz.1",
func(context.Context, json.RawMessage) (any, error) { return "anchor told", nil }, a.respond, nil)
recent, _ := first.Recent()
finished := waitKept(t, keeper, recent[0].ID, CallAnswered)
// A second call, still running when its controller stops.
left := first.begin("mesh-controller", "plans", json.RawMessage(`{}`), "")
waitKept(t, keeper, left.ID, CallRunning)
var said bytes.Buffer
second := NewCallLog()
if err := second.Durably(t.Context(), keeper, "controller@two", log.New(&said, "", 0)); err != nil {
t.Fatal(err)
}
c, found, err := second.Get(finished.ID)
if err != nil || !found {
t.Fatalf("the next controller cannot read %s: %v %v", finished.ID, found, err)
}
if !strings.Contains(string(c.Answer), "anchor told") || c.Caller != "node-tools.g14" || c.Holder != "controller@one" {
t.Fatalf("kept %+v", c)
}
if strings.Contains(string(c.Args), "secret") {
t.Fatalf("a setting was kept: %s", c.Args)
}
abandoned, _, _ := second.Get(left.ID)
if abandoned.State != CallAbandoned || !strings.Contains(said.String(), left.ID) {
t.Fatalf("a call left running reads %q, and was said: %q", abandoned.State, said.String())
}
listed, err := second.Recent()
if err != nil || len(listed) != 2 {
t.Fatalf("the next controller lists %d calls: %v", len(listed), err)
}
// Its own running calls are not its predecessor's: a controller starting again under the same
// name leaves them alone.
mine := second.begin("mesh-controller", "status", nil, "")
waitKept(t, keeper, mine.ID, CallRunning)
if err := second.Durably(t.Context(), keeper, "controller@two", nil); err != nil {
t.Fatal(err)
}
if c, _, _ := keeper.Kept(t.Context(), mine.ID); c.State != CallRunning {
t.Fatalf("a controller marked its own running call %q", c.State)
}
}
// The bucket holds the last thousand calls or fourteen days: a call is two keys, so the stream under
// it holds twice as many messages, and asserting it again keeps the count.
func TestNatsTheCallsBucketIsBounded(t *testing.T) {
callsBucket(t)
js := aBus(t)
if err := js.EnsureControllerBuckets(); err != nil {
t.Fatal(err)
}
info, err := js.Context().StreamInfo("KV_" + broker.CallsBucket)
if err != nil {
t.Fatal(err)
}
if info.Config.MaxMsgs != 2*broker.KeptCallsDurably || info.Config.MaxAge != broker.CallsKeptFor {
t.Fatalf("kept %d messages for %s", info.Config.MaxMsgs, info.Config.MaxAge)
}
}
// An answer larger than the bound is cut and says so, and the record is still written.
func TestNatsAnAnswerLargerThanTheBoundIsCut(t *testing.T) {
keeper := callsBucket(t)
big, _ := json.Marshal(map[string]string{"result": strings.Repeat("x", broker.CallAnswerBytes+10)})
c := Call{ID: "call-1-1", Seat: "mesh-controller", Verb: "plan", State: CallAnswered, Answer: big}
if err := keeper.Keep(t.Context(), c); err != nil {
t.Fatal(err)
}
got, found, err := keeper.Kept(t.Context(), "call-1-1")
if err != nil || !found || !strings.Contains(string(got.Answer), "the first") {
t.Fatalf("%v %v %.200s", found, err, got.Answer)
}
}
// Who asked is read from the inbox the answer goes to.
func TestTheCallerIsTheInboxsPrincipal(t *testing.T) {
for reply, want := range map[string]string{
"_INBOX.node-tools.g14.Pe3sGzAtv8jUBYQ6sKSvKz.1": "node-tools.g14",
"_INBOX.jochen.Pe3sGzAtv8jUBYQ6sKSvKz": "jochen",
"_INBOX.x.1": "",
"mesh.control.one": "",
"": "",
} {
if got := callerOf(reply); got != want {
t.Errorf("%q: %q, want %q", reply, got, want)
}
}
}
+11 -6
View File
@@ -62,7 +62,7 @@ func TestACallThatFinishesInTimeAnswersInFull(t *testing.T) {
if got := a.only()["result"]; got != "all well" {
t.Fatalf("answered %v", got)
}
recent := l.Recent()
recent, _ := l.Recent()
if len(recent) != 1 || recent[0].State != CallAnswered || string(recent[0].Args) != "{}" {
t.Fatalf("kept %+v", recent)
}
@@ -90,12 +90,12 @@ func TestACallThatOutlastsTheWindowSaysItIsRunningAndKeepsItsAnswer(t *testing.T
if got["running"] != true || id == "" || !strings.Contains(got["output"].(string), id) {
t.Fatalf("the running answer does not name its call: %v", got)
}
if c, _ := l.Get(id); c.State != CallRunning {
if c, _, _ := l.Get(id); c.State != CallRunning {
t.Fatalf("while it runs it is kept as %q", c.State)
}
close(release)
<-finished
c, ok := l.Get(id)
c, ok, _ := l.Get(id)
if !ok || c.State != CallFinishedAfter || !strings.Contains(string(c.Answer), "anchor told") {
t.Fatalf("its answer was not kept: %+v", c)
}
@@ -125,7 +125,7 @@ func TestAnAcknowledgedCallIsAnsweredBeforeItGoesOn(t *testing.T) {
if got := a.only()["result"].(map[string]any); got["running"] != true {
t.Fatalf("an acknowledged call answered %v", got)
}
if c := l.Recent()[0]; c.State != CallFinishedAfter || !strings.Contains(string(c.Answer), "sent") {
if c := recentOf(l)[0]; c.State != CallFinishedAfter || !strings.Contains(string(c.Answer), "sent") {
t.Fatalf("kept %+v", c)
}
}
@@ -143,7 +143,7 @@ func TestARefusedAnswerIsKeptAgainstItsCall(t *testing.T) {
if !l.Refusal(refusal, log.New(&logged, "", 0)) {
t.Fatal("the refusal of a kept call's answer was not recognised")
}
c := l.Recent()[0]
c := recentOf(l)[0]
if c.Refused == "" || !strings.Contains(logged.String(), c.ID) {
t.Fatalf("the refusal is not kept or not said: %+v / %q", c, logged.String())
}
@@ -167,8 +167,13 @@ func TestTheLogKeepsTheNewest(t *testing.T) {
for i := 0; i < KeptCalls+5; i++ {
l.begin("s", "v", nil, "")
}
recent := l.Recent()
recent, _ := l.Recent()
if len(recent) != KeptCalls || !strings.HasSuffix(recent[0].ID, "-105") || !strings.HasSuffix(recent[KeptCalls-1].ID, "-6") {
t.Fatalf("kept %d, newest %s, oldest %s", len(recent), recent[0].ID, recent[len(recent)-1].ID)
}
}
func recentOf(l *CallLog) []Call {
out, _ := l.Recent()
return out
}
+9
View File
@@ -380,6 +380,10 @@ func (e Enrolment) Heard(ctx context.Context, report Report) (news bool, err err
if report.Applied == nil && report.Refused == "" && len(report.Failed) == 0 {
if report.Superseded != "" {
log.Printf("%s set aside declaration %s for the newer %s", report.Node, report.Declared, report.Superseded)
} else if err := e.Inventory.RecordHeartbeatGap(ctx, report.Node, node.LastSeen); err != nil {
// The silence before this word, which the machine-silent bound will be set from (novox/hq
// to-be 45 Phase 0). A measurement lost is said and costs the report nothing.
log.Printf("the silence before %s's word could not be recorded: %v", report.Node, err)
}
return false, e.Inventory.Seen(ctx, node.ID)
}
@@ -435,6 +439,11 @@ func (e Enrolment) Heard(ctx context.Context, report Report) (news bool, err err
if err != nil {
return false, err
}
// And how long the machine took, from the send to this first account of it (novox/hq to-be 45
// Phase 0): what the sent-not-reported bound will be set from. Said if lost, never a failure.
if err := e.Inventory.RecordApplyDuration(ctx, report.Node, report.Declared, doing.Outcome); err != nil {
log.Printf("how long %s took to apply could not be recorded: %v", report.Node, err)
}
// A refusal, a failure, or a bare word that the node is there — none of them is an account of
// what the machine holds, so each moves last_seen and nothing else. Recording a partial list
+162
View File
@@ -0,0 +1,162 @@
package link
import (
"context"
"encoding/json"
"errors"
"fmt"
"os"
"sort"
"strconv"
"strings"
"sync/atomic"
"time"
"github.com/nats-io/nats.go"
"github.com/nats-io/nats.go/jetstream"
"github.com/novox/mesh-controller/internal/broker"
)
// The hand-act log (novox/hq to-be 45 §7).
//
// **A repair a person makes by hand is the record of a healer the mesh does not have yet.** Tonight's
// pushes of one machine, plans closed because a report never came, a consumer re-made because it fell
// a week behind — each was done, and the only trace was a chat. So every verb that repairs by hand
// asks why, and writes one entry: who, which verb and arguments, why, when, the condition it
// addresses if one is named, and a cause — a word that, recorded twice in a fortnight, says a healer
// is wanted (S15, from Phase 3). `hand-act record` is the same entry for an act done outside the mesh.
// The controller is the bucket's only writer; its verbs are the way in.
// HandAct is one entry.
type HandAct struct {
ID string `json:"id"`
At time.Time `json:"at"`
// By is who: the bus principal a seat call came from, or the account and machine at a shell.
By string `json:"by"`
// Verb and Args are the act as given: `push`, `plans close`, `broker consumer-reset`, or
// `hand-act record` with what was done outside the mesh.
Verb string `json:"verb"`
Args []string `json:"arguments,omitempty"`
Why string `json:"why"`
// Cause is the condition kind, or a word the person gives; the verb's own name when neither.
Cause string `json:"cause"`
// Condition is the condition's key the act addresses, when it names one.
Condition string `json:"condition,omitempty"`
}
// CallerVar carries a seat call's caller to the command the controller runs for it, so an act done
// through the console says who asked rather than "the controller".
const CallerVar = "MESH_CALLER"
// Caller is who is acting in this process: the seat call's caller when the controller ran it for
// one, otherwise the account and machine at the shell.
func Caller() string {
if c := strings.TrimSpace(os.Getenv(CallerVar)); c != "" {
return c
}
user := os.Getenv("USER")
if user == "" {
user = "an unnamed account"
}
host, _ := os.Hostname()
return fmt.Sprintf("%s at a shell on %s", user, host)
}
type callerKey struct{}
// CallerIn is the caller of the seat call ctx belongs to, empty outside one.
func CallerIn(ctx context.Context) string {
c, _ := ctx.Value(callerKey{}).(string)
return c
}
var handActSeq atomic.Uint64
// RecordHandAct writes one entry. Its key is its time and a sequence, so the bucket lists in order.
func RecordHandAct(ctx context.Context, conn *nats.Conn, act HandAct) (HandAct, error) {
if strings.TrimSpace(act.Why) == "" {
return act, errors.New("an act by hand says why: --why <text>")
}
if act.At.IsZero() {
act.At = time.Now().UTC()
}
if act.ID == "" {
act.ID = "act-" + strconv.FormatInt(act.At.UnixNano(), 10) + "-" + strconv.FormatUint(handActSeq.Add(1), 10)
}
if act.By == "" {
act.By = Caller()
}
if act.Cause == "" {
act.Cause = act.Verb
}
kv, err := handActs(ctx, conn)
if err != nil {
return act, err
}
body, err := json.Marshal(act)
if err != nil {
return act, err
}
_, err = kv.Put(ctx, act.ID, body)
return act, err
}
// HandActs is every entry since a moment, oldest first.
func HandActs(ctx context.Context, conn *nats.Conn, since time.Time) ([]HandAct, error) {
kv, err := handActs(ctx, conn)
if err != nil {
return nil, err
}
w, err := kv.WatchAll(ctx, jetstream.IgnoreDeletes())
if err != nil {
return nil, err
}
defer func() { _ = w.Stop() }()
var out []HandAct
for {
select {
case <-ctx.Done():
return nil, fmt.Errorf("reading the hand-act log: %w", ctx.Err())
case entry := <-w.Updates():
if entry == nil {
sort.SliceStable(out, func(i, j int) bool { return out[i].At.Before(out[j].At) })
return out, nil
}
var a HandAct
if json.Unmarshal(entry.Value(), &a) == nil && !a.At.Before(since) {
out = append(out, a)
}
}
}
}
// RepeatedCauses are the causes recorded more than once within the fortnight before now, with how
// often: each is a repair done by hand again, which is what S15 will raise as a healer wanted.
func RepeatedCauses(acts []HandAct, now time.Time) map[string]int {
counts := map[string]int{}
for _, a := range acts {
if now.Sub(a.At) <= 14*24*time.Hour {
counts[a.Cause]++
}
}
for c, n := range counts {
if n < 2 {
delete(counts, c)
}
}
return counts
}
func handActs(ctx context.Context, conn *nats.Conn) (jetstream.KeyValue, error) {
api, err := jetstream.New(conn)
if err != nil {
return nil, err
}
kv, err := api.KeyValue(ctx, broker.HandActsBucket)
if err != nil {
return nil, fmt.Errorf("the hand-act log %s is not on the bus — the controller asserts it at its "+
"start, so one older than this has not: %w", broker.HandActsBucket, err)
}
return kv, nil
}
+59
View File
@@ -0,0 +1,59 @@
package link
import (
"strings"
"testing"
"time"
"github.com/nats-io/nats.go/jetstream"
"github.com/novox/mesh-controller/internal/broker"
)
// The hand-act log against a real server (novox/hq to-be 45 §7): an act is written with who, why and
// its cause, read back in order, and a cause recorded twice within a fortnight is found.
func TestNatsAnActByHandIsKeptWithWhyAndARepeatIsFound(t *testing.T) {
js := aBus(t)
api, err := jetstream.New(js.Conn())
if err != nil {
t.Fatal(err)
}
_ = api.DeleteKeyValue(t.Context(), broker.HandActsBucket)
if err := js.EnsureControllerBuckets(); err != nil {
t.Fatal(err)
}
t.Setenv(CallerVar, "node-tools.g14, through the mesh-controller seat")
if _, err := RecordHandAct(t.Context(), js.Conn(), HandAct{Verb: "push", Args: []string{"anchor"}}); err == nil {
t.Fatal("an act without why was written")
}
first, err := RecordHandAct(t.Context(), js.Conn(), HandAct{Verb: "push", Args: []string{"anchor"},
Why: "it never reported the send", Cause: "sent-not-reported"})
if err != nil {
t.Fatal(err)
}
if first.By != "node-tools.g14, through the mesh-controller seat" || first.ID == "" {
t.Fatalf("written as %+v", first)
}
if _, err := RecordHandAct(t.Context(), js.Conn(), HandAct{Verb: "plans close", Args: []string{"plan-1"},
Why: "waiting on the same report"}); err != nil {
t.Fatal(err)
}
if _, err := RecordHandAct(t.Context(), js.Conn(), HandAct{Verb: "push", Args: []string{"ace"},
Why: "again", Cause: "sent-not-reported", At: time.Now().UTC().Add(time.Second)}); err != nil {
t.Fatal(err)
}
acts, err := HandActs(t.Context(), js.Conn(), time.Now().Add(-time.Hour))
if err != nil || len(acts) != 3 || acts[0].ID != first.ID || acts[1].Cause != "plans close" {
t.Fatalf("%v %+v", err, acts)
}
repeated := RepeatedCauses(acts, time.Now())
if len(repeated) != 1 || repeated["sent-not-reported"] != 2 {
t.Fatalf("repeated %v", repeated)
}
if old, _ := HandActs(t.Context(), js.Conn(), time.Now().Add(time.Hour)); len(old) != 0 {
t.Fatalf("acts before the moment asked were listed: %+v", old)
}
if !strings.HasPrefix(acts[2].ID, "act-") {
t.Fatalf("%q", acts[2].ID)
}
}