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.
116 lines
3.5 KiB
Go
116 lines
3.5 KiB
Go
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)
|
|
}
|
|
}
|
|
}
|
|
}
|