mesh/merge-gate pass: builds build-agent, mesh-controller, route-proxy → ace, g14, novox, shanks; no bus step; every machine composes with the change as it…
mesh/repo-check fail: its merge-check.sh failed: --- FAIL: TestTheInstallersFirstUserListIsWhatTheControllerWouldCompose (0.79s)
mesh/delivery-group group feat/the-controller-asks-the-operator rejected: a member's own check failed
mesh/delivery superseded: a newer head of the same pull request
With no router, or an ask the router refused and nothing changed since, the controller asked nothing and said it only in its own log. It now keeps a condition of its own, asks-undelivered, naming the conditions not asked and why, cleared once each can be asked again.
299 lines
8.7 KiB
Go
299 lines
8.7 KiB
Go
package main
|
|
|
|
// The asker on the bus: its asks in the controller's bucket `asked`, its asks and cancels published on the
|
|
// seat under the controller's name, the verbs a warrant chooses called with the controller's grant, and the
|
|
// router's record of its asks read under its name (novox/hq ADR 0259).
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"sort"
|
|
"strings"
|
|
"time"
|
|
|
|
"github.com/nats-io/nats.go"
|
|
"github.com/nats-io/nats.go/jetstream"
|
|
|
|
"git.novox.be/novox/mesh-sdk/go/asks"
|
|
|
|
"github.com/novox/mesh-controller/internal/broker"
|
|
"github.com/novox/mesh-controller/internal/catalogue"
|
|
"github.com/novox/mesh-controller/internal/conditions"
|
|
"github.com/novox/mesh-controller/internal/inventory"
|
|
"github.com/novox/mesh-controller/internal/link"
|
|
)
|
|
|
|
// askerFrom is the serving controller's asker; nil in any other process.
|
|
var askerFrom *asker
|
|
|
|
// askWithin is how long a verb a warrant chose is given to answer.
|
|
const askWithin = time.Minute
|
|
|
|
type busAsked struct{ conn *nats.Conn }
|
|
|
|
func (b busAsked) kv(ctx context.Context) (jetstream.KeyValue, error) {
|
|
js, err := jetstream.New(b.conn)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
return js.KeyValue(ctx, broker.AskedBucket)
|
|
}
|
|
|
|
func (b busAsked) Get(ctx context.Context, id string) (*asked, error) {
|
|
kv, err := b.kv(ctx)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
e, err := kv.Get(ctx, id)
|
|
if errors.Is(err, jetstream.ErrKeyNotFound) {
|
|
return nil, nil
|
|
}
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
var r asked
|
|
return &r, json.Unmarshal(e.Value(), &r)
|
|
}
|
|
|
|
func (b busAsked) Put(ctx context.Context, r asked) error {
|
|
kv, err := b.kv(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
body, err := json.Marshal(r)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
_, err = kv.Put(ctx, r.ID, body)
|
|
return err
|
|
}
|
|
|
|
func (b busAsked) All(ctx context.Context) ([]asked, error) {
|
|
kv, err := b.kv(ctx)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
lister, err := kv.ListKeys(ctx)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer func() { _ = lister.Stop() }()
|
|
var out []asked
|
|
for k := range lister.Keys() {
|
|
e, err := kv.Get(ctx, k)
|
|
if err != nil {
|
|
continue
|
|
}
|
|
var r asked
|
|
if json.Unmarshal(e.Value(), &r) == nil {
|
|
out = append(out, r)
|
|
}
|
|
}
|
|
return out, nil
|
|
}
|
|
|
|
// Claim marks an open ask acting, by compare-and-set on its key's revision: only the write that stands acts.
|
|
func (b busAsked) Claim(ctx context.Context, id string, w asks.Warrant) (bool, error) {
|
|
kv, err := b.kv(ctx)
|
|
if err != nil {
|
|
return false, err
|
|
}
|
|
e, err := kv.Get(ctx, id)
|
|
if errors.Is(err, jetstream.ErrKeyNotFound) {
|
|
return false, nil
|
|
}
|
|
if err != nil {
|
|
return false, err
|
|
}
|
|
var r asked
|
|
if err := json.Unmarshal(e.Value(), &r); err != nil {
|
|
return false, err
|
|
}
|
|
if r.State != askOpen || r.Acted != "" {
|
|
return false, nil
|
|
}
|
|
r.State, r.Warrant, r.Acted = string(asks.OutcomeChosen), &w, "acting"
|
|
body, err := json.Marshal(r)
|
|
if err != nil {
|
|
return false, err
|
|
}
|
|
if _, err := kv.Update(ctx, id, body, e.Revision()); err != nil {
|
|
var api *jetstream.APIError
|
|
if errors.Is(err, jetstream.ErrKeyExists) || (errors.As(err, &api) && api.ErrorCode == jetstream.JSErrCodeStreamWrongLastSequence) {
|
|
return false, nil
|
|
}
|
|
return false, err
|
|
}
|
|
return true, nil
|
|
}
|
|
|
|
// callAction performs an action's verb as the controller, through the grant that names it.
|
|
func callAction(conn *nats.Conn) func(ctx context.Context, a conditions.Action, args map[string]string) error {
|
|
return func(ctx context.Context, a conditions.Action, args map[string]string) error {
|
|
seat, verb, ok := strings.Cut(a.Verb, ".")
|
|
if !ok {
|
|
return fmt.Errorf("%q names no seat and verb", a.Verb)
|
|
}
|
|
body := map[string]any{}
|
|
for k, v := range args {
|
|
body[k] = v
|
|
}
|
|
if seat == catalogue.DeliverySeat {
|
|
_, err := askDeliveryOwner(ctx, conn, verb, body)
|
|
return err
|
|
}
|
|
granted := false
|
|
for _, v := range broker.VerbsTheControllerActsOnAWarrant {
|
|
granted = granted || (v.Seat == seat && v.Verb == verb)
|
|
}
|
|
if !granted {
|
|
return fmt.Errorf("%s: %w", a.Verb, errNotGranted)
|
|
}
|
|
var answer link.Answer
|
|
var err error
|
|
if a.Machine != "" {
|
|
answer, err = link.AskSeatTool(ctx, conn, seat, verb, a.Machine, body, askWithin)
|
|
} else {
|
|
answer, err = link.AskMeshSeatTool(ctx, conn, seat, verb, body, askWithin)
|
|
}
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if answer.Error != "" {
|
|
return fmt.Errorf("%s refused: %s", a.Verb, answer.Error)
|
|
}
|
|
return nil
|
|
}
|
|
}
|
|
|
|
// routerRecordOf reads the router's record of one of the controller's asks, under its name, and answers
|
|
// how it ended when it did: the bucket is the one the asks seat's declarer names as its records.
|
|
func routerRecordOf(conn *nats.Conn, inv *inventory.Inventory) func(ctx context.Context, id string) (*asks.Warrant, error) {
|
|
return func(ctx context.Context, id string) (*asks.Warrant, error) {
|
|
bucket, err := asksRecords(ctx, inv)
|
|
if err != nil || bucket == "" {
|
|
return nil, err
|
|
}
|
|
reply, err := conn.RequestWithContext(ctx, "$JS.API.DIRECT.GET.KV_"+bucket+".$KV."+bucket+"."+askerName+"."+id, nil)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
if reply.Header.Get("Status") != "" {
|
|
return nil, nil // none, or not readable: the event says it
|
|
}
|
|
var rec struct {
|
|
State string `json:"state"`
|
|
Warrant *asks.Warrant `json:"warrant"`
|
|
}
|
|
if json.Unmarshal(reply.Data, &rec) != nil || rec.State == "open" || rec.Warrant == nil {
|
|
return nil, nil
|
|
}
|
|
return rec.Warrant, nil
|
|
}
|
|
}
|
|
|
|
// asksRecords is the bucket the asks seat's declarer keeps its record of asks in.
|
|
func asksRecords(ctx context.Context, inv *inventory.Inventory) (string, error) {
|
|
declared, err := inv.Catalogue(ctx)
|
|
if err != nil {
|
|
return "", err
|
|
}
|
|
for _, m := range declared {
|
|
for _, s := range m.DefinesSeats {
|
|
if s.Name == broker.AsksSeat && len(s.Records) > 0 {
|
|
return broker.BucketName(m.Module, s.Records[0]), nil
|
|
}
|
|
}
|
|
}
|
|
return "", nil
|
|
}
|
|
|
|
// routerHereIn says whether a module declaring the asks seat, with its ask named by its caller, is assigned:
|
|
// without it nothing takes an ask, and asking would only fill a queue nobody reads.
|
|
func routerHereIn(inv *inventory.Inventory) func(ctx context.Context) (bool, error) {
|
|
return func(ctx context.Context) (bool, error) {
|
|
entries, err := inv.Catalogued(ctx)
|
|
if err != nil {
|
|
return false, err
|
|
}
|
|
for _, e := range entries {
|
|
for _, s := range e.Manifest.DefinesSeats {
|
|
if s.Name == broker.AsksSeat && s.NamedByCaller("ask") && len(e.On) > 0 {
|
|
return true, nil
|
|
}
|
|
}
|
|
}
|
|
return false, nil
|
|
}
|
|
}
|
|
|
|
// channelsIn is what the channels are now, as a fingerprint: each module claiming a kind of the channel
|
|
// bench, where, promising what, and whether of its own account. An ask the router refused is asked again
|
|
// once this changes.
|
|
func channelsIn(inv *inventory.Inventory) func(ctx context.Context) string {
|
|
return func(ctx context.Context) string {
|
|
entries, err := inv.Catalogued(ctx)
|
|
if err != nil {
|
|
return ""
|
|
}
|
|
var parts []string
|
|
for _, e := range entries {
|
|
for _, c := range e.Manifest.Claims {
|
|
if c.Kind == "" || !catalogue.KindedBenches[c.Name] {
|
|
continue
|
|
}
|
|
on := append([]string(nil), e.On...)
|
|
sort.Strings(on)
|
|
caps := append([]string(nil), c.Capabilities...)
|
|
sort.Strings(caps)
|
|
parts = append(parts, fmt.Sprintf("%s/%s=%s@%s[%s]own:%t", c.Name, c.Kind, e.Manifest.Module,
|
|
strings.Join(on, ","), strings.Join(caps, ","), e.Manifest.RunsAs != ""))
|
|
}
|
|
}
|
|
sort.Strings(parts)
|
|
return strings.Join(parts, ";")
|
|
}
|
|
}
|
|
|
|
// startAsking makes the serving controller's asker and hands it the router's words.
|
|
func startAsking(ctx context.Context, open *stores, server *link.Server, conn *nats.Conn, keeper *conditions.Keeper) {
|
|
js, err := jetstream.New(conn)
|
|
if err != nil {
|
|
fmt.Printf("the operator cannot be asked: %v\n", err)
|
|
return
|
|
}
|
|
a := &asker{
|
|
open: keeper.Open,
|
|
silence: func(ctx context.Context, key string, d time.Duration, by, why string) error {
|
|
_, err := keeper.Silence(ctx, key, d, by, why)
|
|
return err
|
|
},
|
|
store: busAsked{conn: conn},
|
|
publish: func(ctx context.Context, subject string, body []byte, id string) error {
|
|
_, err := js.Publish(ctx, subject, body, jetstream.WithMsgID(id))
|
|
return err
|
|
},
|
|
call: callAction(conn),
|
|
record: func(ctx context.Context, act link.HandAct) error {
|
|
_, err := link.RecordHandAct(ctx, conn, act)
|
|
return err
|
|
},
|
|
routerRecord: routerRecordOf(conn, open.inventory),
|
|
routerHere: routerHereIn(open.inventory),
|
|
channels: channelsIn(open.inventory),
|
|
raise: func(ctx context.Context, obs []conditions.Observation) error {
|
|
return keeper.Reconcile(ctx, sourceAsker, obs)
|
|
},
|
|
now: time.Now,
|
|
logf: func(format string, args ...any) { fmt.Printf(format+"\n", args...) },
|
|
}
|
|
if err := server.Decides(a); err != nil {
|
|
fmt.Printf("the operator's answers cannot be heard, so nothing is asked: %v\n", err)
|
|
return
|
|
}
|
|
askerFrom = a
|
|
go a.keep(ctx)
|
|
}
|