Prove the operator's answers end to end before anybody uses them (hq ADR 0259) #74

Merged
mesh-admin merged 9 commits from proofs/asks-answered-on-the-phone into main 2026-10-10 12:15:06 +00:00
10 changed files with 1538 additions and 3 deletions
+207
View File
@@ -0,0 +1,207 @@
package asks
import (
"encoding/json"
"sync"
"testing"
"time"
stdio "git.novox.be/novox/mesh-sdk/go"
"git.novox.be/novox/mesh-sdk/go/asker"
"git.novox.be/novox/mesh-sdk/go/asks"
"github.com/nats-io/nats.go"
)
// labAsker is the module `lab-asker` asking the operator, as any module does: the SDK's asker client
// (mesh-sdk go/asker) over the lab's bus, on that module's own composed credential. What it performs on a
// warrant is one switch flipped, bound when it asked, and counted: the proof that an approval does exactly
// that one thing, once.
type labAsker struct {
t *testing.T
client *asker.Client
nc *nats.Conn
mu sync.Mutex
warrants []asks.Warrant
handled []asker.Handled
acts []map[string]string // what was performed
}
// flip is the act an approve option binds.
func flip(which string) asks.Act {
return asks.Act{"verb": "lab.flip", "arg.switch": which}
}
func newLabAsker(t *testing.T, l *lab) *labAsker {
nc, _ := l.as(machine + ".lab-asker")
js, err := nc.JetStream()
if err != nil {
t.Fatal(err)
}
a := &labAsker{t: t, nc: nc}
a.client = &asker.Client{Module: "lab-asker", State: &memState{vals: map[string]json.RawMessage{}, revs: map[string]uint64{}},
Now: time.Now,
Publish: func(subject string, body any, id string) error {
raw, err := json.Marshal(body)
if err != nil {
return err
}
_, err = js.Publish(subject, raw, nats.MsgId("lab-asker."+id))
return err
}}
// Its warrants, on the one subject the bus lets it hear them on.
sub, err := nc.Subscribe(asks.DecidedSubject("lab-asker"), func(m *nats.Msg) {
var w asks.Warrant
if err := json.Unmarshal(m.Data, &w); err != nil {
t.Errorf("a word on the asker's decided subject is not a warrant: %v", err)
return
}
a.take(w)
})
if err != nil {
t.Fatal(err)
}
t.Cleanup(func() { _ = sub.Unsubscribe() })
return a
}
// take hands one warrant to the SDK's Take, which performs the option once, and only the act it bound.
func (a *labAsker) take(w asks.Warrant) asker.Handled {
h, err := a.client.Take(w, func(o asks.Option, w asks.Warrant) error {
if o.Level == asks.Acknowledge {
return nil // decline: nothing is performed
}
act := flip(o.ID)
if err := o.Performs(act); err != nil {
return err
}
a.mu.Lock()
a.acts = append(a.acts, act)
a.mu.Unlock()
return nil
})
a.mu.Lock()
a.warrants = append(a.warrants, w)
a.handled = append(a.handled, h)
a.mu.Unlock()
if err != nil {
a.t.Logf("the asker refused a warrant for %s: %v", w.Ask, err)
}
return h
}
// ask asks the operator: one approve option binding a switch, one decline that acknowledges.
func (a *labAsker) ask(id, headline string, expires time.Duration) asks.Ask {
a.t.Helper()
binds, err := asks.ActDigest(flip("approve"))
if err != nil {
a.t.Fatal(err)
}
q := asks.Ask{ID: id, Headline: headline,
Explanation: "Needs you: approve or decline. The lab asks whether to flip one switch.",
Who: asks.Operator, Expires: time.Now().Add(expires), OnExpiry: "nothing is flipped",
Options: []asks.Option{
{ID: "approve", Label: "Approve", Does: "the switch is flipped, once", Level: asks.Approve, Binds: binds},
{ID: "decline", Label: "Decline", Does: "nothing is done", Level: asks.Acknowledge},
}}
if err := a.client.Ask(q); err != nil {
a.t.Fatalf("the asker could not ask %s: %v", id, err)
}
return q
}
// warrantFor waits for the router's word on one ask.
func (a *labAsker) warrantFor(id string, within time.Duration) asks.Warrant {
a.t.Helper()
deadline := time.Now().Add(within)
for time.Now().Before(deadline) {
a.mu.Lock()
for _, w := range a.warrants {
if w.Ask == id {
a.mu.Unlock()
return w
}
}
a.mu.Unlock()
time.Sleep(100 * time.Millisecond)
}
a.t.Fatalf("no word on the ask %s within %s", id, within)
return asks.Warrant{}
}
// noWordOn says the router said nothing on an ask for a while.
func (a *labAsker) noWordOn(id string, wait time.Duration) bool {
time.Sleep(wait)
a.mu.Lock()
defer a.mu.Unlock()
for _, w := range a.warrants {
if w.Ask == id {
return false
}
}
return true
}
func (a *labAsker) performed() []map[string]string {
a.mu.Lock()
defer a.mu.Unlock()
return append([]map[string]string(nil), a.acts...)
}
func (a *labAsker) wordsOn(id string) int {
a.mu.Lock()
defer a.mu.Unlock()
n := 0
for _, w := range a.warrants {
if w.Ask == id {
n++
}
}
return n
}
// memState is the asker's own record, kept in memory: the lab's asker restarts never.
type memState struct {
mu sync.Mutex
vals map[string]json.RawMessage
revs map[string]uint64
n uint64
}
func (m *memState) Get(k string) (*stdio.StateEntry, error) {
m.mu.Lock()
defer m.mu.Unlock()
v, ok := m.vals[k]
if !ok {
return nil, nil
}
return &stdio.StateEntry{Key: k, Value: v, Revision: m.revs[k]}, nil
}
func (m *memState) Create(k string, v any) (uint64, error) {
m.mu.Lock()
defer m.mu.Unlock()
if _, ok := m.vals[k]; ok {
return 0, stdio.ErrChanged
}
return m.put(k, v)
}
func (m *memState) Update(k string, v any, rev uint64) (uint64, error) {
m.mu.Lock()
defer m.mu.Unlock()
if rev == 0 || m.revs[k] != rev {
return 0, stdio.ErrChanged
}
return m.put(k, v)
}
func (m *memState) put(k string, v any) (uint64, error) {
raw, err := json.Marshal(v)
if err != nil {
return 0, err
}
m.n++
m.vals[k], m.revs[k] = raw, m.n
return m.n, nil
}
+4
View File
@@ -0,0 +1,4 @@
// Package asks is the lab's proof of the operator's answers (novox/hq ADR 0259): the whole flow, from an ask
// to the act its approval authorises, on a real bus composed by the controller, with a fake Telegram. The
// proof is TestTheOperatorsAnswerEndToEnd; the live acceptance after rollout is live-acceptance.sh.
package asks
+282
View File
@@ -0,0 +1,282 @@
package asks
import (
"encoding/json"
"fmt"
"net/http"
"net/http/httptest"
"strings"
"sync"
"testing"
"time"
)
// fakeTelegram is the Bot API as the channel uses it, and nothing more: the lab's phone. It holds the updates
// the lab makes happen (a message typed, a button tapped) for the channel's long poll, and keeps every message
// the bot sends or edits, by chat, so the lab reads what the operator would see. It never reaches Telegram.
type fakeTelegram struct {
t *testing.T
server *httptest.Server
token string
mu sync.Mutex
updates []map[string]any
next int64
arrived chan struct{}
sent []tgMessage
msgID int64
taps int // answerCallbackQuery calls
deleted int
}
// tgMessage is one message the bot sent or edited: its chat, its id, its words and its buttons (label →
// callback data).
type tgMessage struct {
Chat int64
ID int64
Text string
Buttons [][2]string
Edited bool
At time.Time
}
func newFakeTelegram(t *testing.T, token string) *fakeTelegram {
f := &fakeTelegram{t: t, token: token, next: 1, arrived: make(chan struct{}, 1)}
f.server = httptest.NewServer(http.HandlerFunc(f.serve))
t.Cleanup(f.server.Close)
return f
}
func (f *fakeTelegram) URL() string { return f.server.URL }
func (f *fakeTelegram) serve(w http.ResponseWriter, r *http.Request) {
prefix := "/bot" + f.token + "/"
if !strings.HasPrefix(r.URL.Path, prefix) {
// A wrong token is what the real API answers 401 to.
w.WriteHeader(http.StatusUnauthorized)
_ = json.NewEncoder(w).Encode(map[string]any{"ok": false, "error_code": 401, "description": "Unauthorized"})
return
}
method := strings.TrimPrefix(r.URL.Path, prefix)
var body map[string]any
_ = json.NewDecoder(r.Body).Decode(&body)
ok := func(result any) { _ = json.NewEncoder(w).Encode(map[string]any{"ok": true, "result": result}) }
switch method {
case "getMe":
ok(map[string]any{"id": 1000, "is_bot": true, "first_name": "Lab", "username": "lab_mesh_bot"})
case "getUpdates":
offset := int64(num(body["offset"]))
wait := time.Duration(num(body["timeout"])) * time.Second
if wait > 2*time.Second {
wait = 2 * time.Second // the lab's phone answers quickly; the channel polls again at once
}
deadline := time.Now().Add(wait)
for {
f.mu.Lock()
var out []map[string]any
for _, u := range f.updates {
if int64(num(u["update_id"])) >= offset {
out = append(out, u)
}
}
f.mu.Unlock()
if len(out) > 0 {
f.t.Logf("phone: getUpdates from %d hands %d update(s)", offset, len(out))
}
if len(out) > 0 || time.Now().After(deadline) {
ok(orEmpty(out))
return
}
select {
case <-f.arrived:
case <-time.After(time.Until(deadline)):
}
}
case "sendMessage", "editMessageText":
chat := int64(num(body["chat_id"]))
text, _ := body["text"].(string)
m := tgMessage{Chat: chat, Text: text, Buttons: buttonsOf(body["reply_markup"]), At: time.Now()}
f.mu.Lock()
if method == "sendMessage" {
f.msgID++
m.ID = f.msgID
} else {
m.ID, m.Edited = int64(num(body["message_id"])), true
}
f.sent = append(f.sent, m)
f.mu.Unlock()
f.t.Logf("phone: %s to %d (#%d): %q %v", method, chat, m.ID, firstLine(text), m.Buttons)
ok(map[string]any{"message_id": m.ID, "chat": map[string]any{"id": chat, "type": "private"}, "text": text})
case "answerCallbackQuery":
f.mu.Lock()
f.taps++
f.mu.Unlock()
ok(true)
case "deleteMessage":
f.mu.Lock()
f.deleted++
f.mu.Unlock()
ok(true)
default:
w.WriteHeader(http.StatusNotFound)
_ = json.NewEncoder(w).Encode(map[string]any{"ok": false, "error_code": 404, "description": "Not Found: " + method})
}
}
// push makes one thing happen on the phone, for the channel's next poll.
func (f *fakeTelegram) push(u map[string]any) {
f.mu.Lock()
u["update_id"] = f.next
f.next++
f.updates = append(f.updates, u)
f.mu.Unlock()
select {
case f.arrived <- struct{}{}:
default:
}
}
// typed is a message an account writes in a chat with the bot; private when the chat is the account's own.
func (f *fakeTelegram) typed(from int64, chat int64, chatType, text string) {
f.push(map[string]any{"message": map[string]any{"message_id": 9000 + f.nextID(), "text": text,
"from": map[string]any{"id": from, "is_bot": false, "first_name": fmt.Sprintf("Account %d", from)},
"chat": map[string]any{"id": chat, "type": chatType}}})
}
// tapped is a button tapped on a message the bot sent, carrying the button's data and the message's words as
// the service gives them: the words the message shows now (its last send or edit), as on the phone.
func (f *fakeTelegram) tapped(from int64, chat int64, chatType string, onMessage int64, data string) {
f.tappedOn(from, chat, chatType, onMessage, data, f.textOf(onMessage))
}
// tappedOn is a tap on a message showing the words given: what a message changed under the operator's thumb,
// or forwarded and edited elsewhere, carries.
func (f *fakeTelegram) tappedOn(from int64, chat int64, chatType string, onMessage int64, data, text string) {
f.push(map[string]any{"callback_query": map[string]any{"id": fmt.Sprintf("cb-%d", f.nextID()), "data": data,
"from": map[string]any{"id": from, "is_bot": false, "first_name": fmt.Sprintf("Account %d", from)},
"message": map[string]any{"message_id": onMessage, "text": text,
"chat": map[string]any{"id": chat, "type": chatType}}}})
}
// textOf is what a message the bot sent shows now.
func (f *fakeTelegram) textOf(id int64) string {
f.mu.Lock()
defer f.mu.Unlock()
text := ""
for _, m := range f.sent {
if m.ID == id {
text = m.Text
}
}
return text
}
func (f *fakeTelegram) nextID() int64 {
f.mu.Lock()
defer f.mu.Unlock()
f.msgID++
return f.msgID
}
// waitFor waits for a message to a chat that satisfies match, sent after since.
func (f *fakeTelegram) waitFor(what string, chat int64, since time.Time, within time.Duration, match func(tgMessage) bool) tgMessage {
f.t.Helper()
deadline := time.Now().Add(within)
for time.Now().Before(deadline) {
f.mu.Lock()
for _, m := range f.sent {
if m.Chat == chat && !m.At.Before(since) && match(m) {
f.mu.Unlock()
return m
}
}
f.mu.Unlock()
time.Sleep(100 * time.Millisecond)
}
f.t.Fatalf("the phone never showed %s (chat %d) within %s; it showed:\n%s", what, chat, within, f.transcript())
return tgMessage{}
}
// since is every message to a chat sent at or after a time.
func (f *fakeTelegram) since(chat int64, at time.Time) []tgMessage {
f.mu.Lock()
defer f.mu.Unlock()
var out []tgMessage
for _, m := range f.sent {
if m.Chat == chat && !m.At.Before(at) {
out = append(out, m)
}
}
return out
}
func (f *fakeTelegram) transcript() string {
f.mu.Lock()
defer f.mu.Unlock()
var b strings.Builder
for _, m := range f.sent {
fmt.Fprintf(&b, " %s to %d #%d edited=%v %q %v\n", m.At.Format("15:04:05.000"), m.Chat, m.ID, m.Edited,
firstLine(m.Text), m.Buttons)
}
return b.String()
}
// button is the data of the button labelled so on a message.
func button(m tgMessage, label string) string {
for _, b := range m.Buttons {
if b[0] == label {
return b[1]
}
}
return ""
}
func labels(m tgMessage) []string {
var out []string
for _, b := range m.Buttons {
out = append(out, b[0])
}
return out
}
func buttonsOf(markup any) [][2]string {
var out [][2]string
m, _ := markup.(map[string]any)
rows, _ := m["inline_keyboard"].([]any)
for _, row := range rows {
cells, _ := row.([]any)
for _, c := range cells {
cell, _ := c.(map[string]any)
text, _ := cell["text"].(string)
data, _ := cell["callback_data"].(string)
out = append(out, [2]string{text, data})
}
}
return out
}
func num(v any) float64 {
switch n := v.(type) {
case float64:
return n
case int64:
return float64(n)
case int:
return float64(n)
}
return 0
}
func orEmpty(in []map[string]any) []map[string]any {
if in == nil {
return []map[string]any{}
}
return in
}
func firstLine(s string) string {
if i := strings.IndexByte(s, '\n'); i >= 0 {
return s[:i] + " …"
}
return s
}
+22
View File
@@ -0,0 +1,22 @@
module github.com/novox/mesh-lab/asks
go 1.26.0
require (
git.novox.be/novox/mesh-sdk/go v0.1.10
github.com/nats-io/nats-server/v2 v2.11.17
github.com/nats-io/nats.go v1.54.0
)
require (
github.com/antithesishq/antithesis-sdk-go v0.7.0-default-no-op // indirect
github.com/google/go-tpm v0.9.8 // indirect
github.com/klauspost/compress v1.20.0 // indirect
github.com/minio/highwayhash v1.0.4 // indirect
github.com/nats-io/jwt/v2 v2.8.1 // indirect
github.com/nats-io/nkeys v0.4.16 // indirect
github.com/nats-io/nuid v1.0.1 // indirect
golang.org/x/crypto v0.57.0 // indirect
golang.org/x/sys v0.48.0 // indirect
golang.org/x/time v0.15.0 // indirect
)
+27
View File
@@ -0,0 +1,27 @@
git.novox.be/novox/mesh-sdk/go v0.1.10 h1:EHhnPRrkw7/CSQAqRsxYnTvlTZ3dI02IPO6WSmVX80o=
git.novox.be/novox/mesh-sdk/go v0.1.10/go.mod h1:GFuZUElBZ9A++mxgIKo97aXXo+kV0uJ/UkbhQPPIbrY=
github.com/antithesishq/antithesis-sdk-go v0.7.0-default-no-op h1:Z/MZK75wC/NSrkgqeNIa7jexam9uWzhLmFTSCPI/kn0=
github.com/antithesishq/antithesis-sdk-go v0.7.0-default-no-op/go.mod h1:FQyySiasQQM8735Ddel3MRojmy4dA1IqCeyJ5jmPMbI=
github.com/google/go-tpm v0.9.8 h1:slArAR9Ft+1ybZu0lBwpSmpwhRXaa85hWtMinMyRAWo=
github.com/google/go-tpm v0.9.8/go.mod h1:h9jEsEECg7gtLis0upRBQU+GhYVH6jMjrFxI8u6bVUY=
github.com/klauspost/compress v1.20.0 h1:a3C1ke2ohxFymNlb2HWAHjDeKCI90scRskErZkR0ezA=
github.com/klauspost/compress v1.20.0/go.mod h1:LUdAzn7YLVvxLpc7y3V1m40wESHTgc1422pwwBSKYuI=
github.com/minio/highwayhash v1.0.4 h1:asJizugGgchQod2ja9NJlGOWq4s7KsAWr5XUc9Clgl4=
github.com/minio/highwayhash v1.0.4/go.mod h1:GGYsuwP/fPD6Y9hMiXuapVvlIUEhFhMTh0rxU3ik1LQ=
github.com/nats-io/jwt/v2 v2.8.1 h1:V0xpGuD/N8Mi+fQNDynXohVvp7ZztevW5io8CUWlPmU=
github.com/nats-io/jwt/v2 v2.8.1/go.mod h1:nWnOEEiVMiKHQpnAy4eXlizVEtSfzacZ1Q43LIRavZg=
github.com/nats-io/nats-server/v2 v2.11.17 h1:GKEghcFK6A+aFx11Yf1LjgLC3txAwvyhnYzhBIQZA8I=
github.com/nats-io/nats-server/v2 v2.11.17/go.mod h1:B1sFVz4StNosQ903ak4N1G01Fl/9f8e06mXpFIE2K24=
github.com/nats-io/nats.go v1.54.0 h1:vsXoOxjHp/GmPUN+EcI7uOf/uB+iAP+kEsAFNQN0yzA=
github.com/nats-io/nats.go v1.54.0/go.mod h1:y+DZoD1oBOYfZTU681eTUiUjI0vbqYGixNVFHcjHJ0k=
github.com/nats-io/nkeys v0.4.16 h1:rd5oAuLOb8mnAycB0xleuEBNS1pVVnN0fv/FF34Eypg=
github.com/nats-io/nkeys v0.4.16/go.mod h1:llLgWoI0o4z/Q57q2R1kHfmocyhGV6VG/U18Glg1Afs=
github.com/nats-io/nuid v1.0.1 h1:5iA8DT8V7q8WK2EScv2padNa/rTESc1KdnPw4TC2paw=
github.com/nats-io/nuid v1.0.1/go.mod h1:19wcPz3Ph3q0Jbyiqsd0kePYG7A95tJPxeL+1OSON2c=
golang.org/x/crypto v0.57.0 h1:3ZVCjf8Ggz7zneR/EHRVx68Ctf+2pmIMP2UFhh9cC6M=
golang.org/x/crypto v0.57.0/go.mod h1:Fdz0i5U6CoizGwLda9DttjSk6qlZo25zYNtR+ycvuZA=
golang.org/x/sys v0.21.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA=
golang.org/x/sys v0.48.0 h1:bbX/i/6MgT9BVLM9RT1thmxL04yeTAhbEz4SyadbXoo=
golang.org/x/sys v0.48.0/go.mod h1:hNLxWAXmnKAxqDtdwIYC4bM9oQPEecfsnNMuSxOs3og=
golang.org/x/time v0.15.0 h1:bbrp8t3bGUeFOx08pvsMYRTCVSMk89u4tKbNOZbp88U=
golang.org/x/time v0.15.0/go.mod h1:Y4YMaQmXwGQZoFaVFk4YpCt4FLQMYKZe9oeV/f4MSno=
+537
View File
@@ -0,0 +1,537 @@
package asks
import (
"bytes"
"context"
"encoding/json"
"fmt"
"net"
"os"
"os/exec"
"path/filepath"
"regexp"
"strings"
"sync"
"testing"
"time"
"github.com/nats-io/nats-server/v2/server"
"github.com/nats-io/nats.go"
"github.com/nats-io/nats.go/jetstream"
)
// machine is the lab's one machine, as the controller's fixture names it.
const machine = "anchor"
// credential is what a lab process connects as (the controller fixture's credentials.json, mesh-tools'
// bus.Credential).
type credential struct {
URL string `json:"url"`
Node string `json:"node,omitempty"`
Module string `json:"module,omitempty"`
User string `json:"user"`
Password string `json:"password"`
}
// lab is one run of the proof: the repositories it builds from, its directory, its bus, its phone.
type lab struct {
t *testing.T
repos clones
dir string
url string
creds map[string]credential
server *server.Server
// observer is the lab's own user on the bus, outside the mesh's composition: it plays the controller's
// conditions verb and reads the router's state, and is never used to show that something is allowed.
observer *nats.Conn
tg *fakeTelegram
mu sync.Mutex
conditions []map[string]any
// notFree are the machines the controller's root-free verb answers as not root-free.
notFree map[string]string
rootAsked int
desk []string // what the desk was shown
procs []*exec.Cmd
logs map[string]*bytes.Buffer
}
// clones is where each core repository the proof builds is checked out: side by side in MESH_LAB_REPOS (beside
// mesh-lab by default), or one by one in MESH_LAB_CONTROLLER, MESH_LAB_TOOLS and MESH_LAB_CATALOGUE — a branch
// under review is often a worktree of its own.
type clones struct{ controller, tools, catalogue string }
func clonesFor(t *testing.T) clones {
repos := os.Getenv("MESH_LAB_REPOS")
if repos == "" {
repos = filepath.Join("..", "..")
}
at := func(env, name string) string {
dir := os.Getenv(env)
if dir == "" {
dir = filepath.Join(repos, name)
}
abs, err := filepath.Abs(dir)
if err != nil {
t.Fatal(err)
}
return abs
}
c := clones{controller: at("MESH_LAB_CONTROLLER", "mesh-controller"), tools: at("MESH_LAB_TOOLS", "mesh-tools"),
catalogue: at("MESH_LAB_CATALOGUE", "mesh-catalog")}
for _, want := range []string{filepath.Join(c.controller, "go.mod"), filepath.Join(c.tools, "node-tools", "go.mod"),
filepath.Join(c.catalogue, "modules", "messenger", "go.mod"), filepath.Join(c.catalogue, "modules", "telegram", "go.mod")} {
if _, err := os.Stat(want); err != nil {
t.Skipf("NOT RUN: the proof builds the controller, the runtime, the router and the Telegram channel from "+
"their clones, and there is no %s (set MESH_LAB_REPOS, or MESH_LAB_CONTROLLER, MESH_LAB_TOOLS and "+
"MESH_LAB_CATALOGUE)", want)
}
}
// The controller beside is, on the build seat, the one the mesh runs: until it carries the proof's fixture
// (mesh-controller's TestTheAsksLabBus), the proof cannot compose its bus, and says so rather than fail
// every pull request of the lab for a merge not yet made.
fixtures, _ := filepath.Glob(filepath.Join(c.controller, "internal", "inventory", "*_test.go"))
has := false
for _, f := range fixtures {
if body, err := os.ReadFile(f); err == nil && strings.Contains(string(body), "func TestTheAsksLabBus(") {
has = true
}
}
if !has {
t.Skipf("NOT RUN: the controller at %s has no TestTheAsksLabBus, the fixture that composes the proof's bus "+
"(mesh-controller #157): the proof runs once the controller it is given carries it", c.controller)
}
return c
}
// run runs a command in a directory and fails the test with what it said.
func (l *lab) run(dir string, env []string, name string, args ...string) string {
l.t.Helper()
cmd := exec.Command(name, args...)
cmd.Dir = dir
cmd.Env = append(os.Environ(), env...)
out, err := cmd.CombinedOutput()
if err != nil {
l.t.Fatalf("%s %s in %s: %v\n%s", name, strings.Join(args, " "), dir, err, out)
}
return string(out)
}
// compose asks the controller at its commit for the bus it would compose, then (bus set) to raise it.
func (l *lab) compose(bus string) {
// The anchor composed as a machine the controller measured root-free: verified-sender reaches the router
// only then (ADR 0259 §8). The not-free case is played live, through the root-free verb, in the proof.
env := []string{"GOFLAGS=-mod=vendor", "GOPROXY=off", "MESH_LAB_ASKS_OUT=" + l.dir, "MESH_LAB_ASKS_ROOT_FREE=true",
"MESH_LAB_ASKS_CATALOGUE=" + filepath.Join(l.repos.catalogue, "modules"),
"MESH_LAB_ASKS_RUNTIME=" + filepath.Join(l.repos.tools, "node-tools", "module.json")}
if bus != "" {
env = append(env, "MESH_LAB_ASKS_BUS="+bus)
}
out := l.run(l.repos.controller, env, "go", "test", "-count=1", "-v", "-run",
"^TestTheAsksLabBus$", "./internal/inventory/")
if !strings.Contains(out, "--- PASS: TestTheAsksLabBus") {
l.t.Fatalf("the controller did not compose the lab's bus (is its fixture on this commit?):\n%s", out)
}
}
// startBus starts a server of the mesh's release with the controller's composed accounts, and the lab's
// observer beside them.
func (l *lab) startBus() {
l.t.Helper()
accounts, err := os.ReadFile(filepath.Join(l.dir, "accounts.conf"))
if err != nil {
l.t.Fatal(err)
}
const usersOpen = " users = [\n"
if !bytes.Contains(accounts, []byte(usersOpen)) {
l.t.Fatalf("the composed accounts block has no users list:\n%s", accounts)
}
observerPassword := fmt.Sprintf("observer-%d", time.Now().UnixNano())
accounts = bytes.Replace(accounts, []byte(usersOpen),
[]byte(usersOpen+fmt.Sprintf(" { user: \"lab-observer\", password: %q }\n", observerPassword)), 1)
port := freePort(l.t)
conf := fmt.Sprintf("listen: \"127.0.0.1:%d\"\njetstream {\n store_dir: %q\n}\n\n%s", port,
filepath.Join(l.dir, "jetstream"), accounts)
path := filepath.Join(l.dir, "bus.conf")
if err := os.WriteFile(path, []byte(conf), 0o600); err != nil {
l.t.Fatal(err)
}
opts, err := server.ProcessConfigFile(path)
if err != nil {
l.t.Fatalf("the composed configuration is not one the server reads: %v", err)
}
// The server's own log, in the lab's directory: a refusal it says is evidence, kept with the run.
opts.NoSigs, opts.LogFile = true, filepath.Join(l.dir, "bus.log")
opts.Trace = os.Getenv("MESH_LAB_ASKS_TRACE") != ""
s, err := server.NewServer(opts)
if err != nil {
l.t.Fatal(err)
}
s.ConfigureLogger()
go s.Start()
if !s.ReadyForConnections(30 * time.Second) {
l.t.Fatal("the lab's bus did not come up")
}
l.server, l.url = s, s.ClientURL()
l.t.Cleanup(func() { s.Shutdown(); s.WaitForShutdown() })
raw, err := os.ReadFile(filepath.Join(l.dir, "credentials.json"))
if err != nil {
l.t.Fatal(err)
}
if err := json.Unmarshal(raw, &l.creds); err != nil {
l.t.Fatal(err)
}
l.observer, err = nats.Connect(l.url, nats.UserInfo("lab-observer", observerPassword))
if err != nil {
l.t.Fatal(err)
}
l.t.Cleanup(l.observer.Close)
}
func freePort(t *testing.T) int {
ln, err := net.Listen("tcp", "127.0.0.1:0")
if err != nil {
t.Fatal(err)
}
defer ln.Close()
return ln.Addr().(*net.TCPAddr).Port
}
// as connects as one of the composed users. refused collects what the server refused it.
func (l *lab) as(user string) (*nats.Conn, *refusals) {
l.t.Helper()
c, ok := l.creds[user]
if !ok {
l.t.Fatalf("the controller composed no user %s", user)
}
r := &refusals{}
nc, err := nats.Connect(l.url, nats.UserInfo(c.User, c.Password), nats.CustomInboxPrefix("_INBOX."+c.User),
nats.ErrorHandler(func(_ *nats.Conn, _ *nats.Subscription, err error) { r.add(err) }))
if err != nil {
l.t.Fatalf("connecting as %s: %v", user, err)
}
l.t.Cleanup(nc.Close)
return nc, r
}
// refusals are the permission violations the server said to one connection.
type refusals struct {
mu sync.Mutex
said []string
}
func (r *refusals) add(err error) {
r.mu.Lock()
defer r.mu.Unlock()
r.said = append(r.said, err.Error())
}
func (r *refusals) of(subject string) bool {
r.mu.Lock()
defer r.mu.Unlock()
for _, s := range r.said {
if strings.Contains(strings.ToLower(s), "permissions violation") && strings.Contains(s, subject) {
return true
}
}
return false
}
// refusedToPublish publishes once as a user and says whether the server refused it.
func (l *lab) refusedToPublish(user, subject string, body []byte) bool {
l.t.Helper()
nc, r := l.as(user)
if err := nc.Publish(subject, body); err != nil {
return true
}
_ = nc.Flush()
deadline := time.Now().Add(2 * time.Second)
for time.Now().Before(deadline) {
if r.of(subject) {
return true
}
time.Sleep(50 * time.Millisecond)
}
return false
}
// build builds one Go program from a directory into the lab's bin.
func (l *lab) build(dir, pkg, name string, ldflags string) string {
out := filepath.Join(l.dir, "bin", name)
args := []string{"build", "-o", out}
if ldflags != "" {
args = append(args, "-ldflags", ldflags)
}
l.run(dir, []string{"GOPRIVATE=git.novox.be", "GOFLAGS=-mod=mod", "CGO_ENABLED=0"}, "go", append(args, pkg)...)
return out
}
// runtime starts the mesh's runtime on a module's own credential, serving that module's bundle alone: how
// the mesh runs a module that `runs-as` an account of its own.
func (l *lab) runtime(module, runtimeBin, bundle string, env map[string]string) {
l.t.Helper()
user := machine + "." + module
c := l.creds[user]
c.URL, c.Node, c.Module = l.url, machine, module
credPath := filepath.Join(l.dir, module, "broker")
if err := os.MkdirAll(filepath.Dir(credPath), 0o700); err != nil {
l.t.Fatal(err)
}
raw, _ := json.Marshal(c)
if err := os.WriteFile(credPath, raw, 0o600); err != nil {
l.t.Fatal(err)
}
toolEnv, _ := json.Marshal(map[string]map[string]string{module: env})
cmd := exec.Command(runtimeBin, "serve")
cmd.Env = append(os.Environ(), "MESH_BROKER_FILE="+credPath, "MESH_TOOL_MODULES="+module+"="+bundle,
"MESH_TOOL_ENV="+string(toolEnv))
logs := &bytes.Buffer{}
cmd.Stdout, cmd.Stderr = &lockedWriter{w: logs, mu: &l.mu}, &lockedWriter{w: logs, mu: &l.mu}
if err := cmd.Start(); err != nil {
l.t.Fatal(err)
}
l.mu.Lock()
l.procs = append(l.procs, cmd)
if l.logs == nil {
l.logs = map[string]*bytes.Buffer{}
}
l.logs[module] = logs
l.mu.Unlock()
l.t.Cleanup(func() {
_ = cmd.Process.Signal(os.Interrupt)
done := make(chan struct{})
go func() { _ = cmd.Wait(); close(done) }()
select {
case <-done:
case <-time.After(5 * time.Second):
_ = cmd.Process.Kill()
}
if l.t.Failed() || os.Getenv("MESH_LAB_ASKS_SAY") != "" {
l.mu.Lock()
l.t.Logf("--- what the %s runtime said:\n%s", module, logs.String())
l.mu.Unlock()
}
})
}
// said is what a runtime said so far.
func (l *lab) said(module string) string {
l.mu.Lock()
defer l.mu.Unlock()
if b := l.logs[module]; b != nil {
return b.String()
}
return ""
}
// waitSaid waits for a runtime to say something.
func (l *lab) waitSaid(module, what string, within time.Duration) {
l.t.Helper()
deadline := time.Now().Add(within)
for time.Now().Before(deadline) {
if strings.Contains(l.said(module), what) {
return
}
time.Sleep(200 * time.Millisecond)
}
l.t.Fatalf("the %s runtime never said %q within %s:\n%s", module, what, within, l.said(module))
}
type lockedWriter struct {
w *bytes.Buffer
mu *sync.Mutex
}
func (w *lockedWriter) Write(p []byte) (int, error) {
w.mu.Lock()
defer w.mu.Unlock()
return w.w.Write(p)
}
// playController answers the controller's `conditions` verb from the lab's list, and its `root-free` verb from
// the lab's word on each machine (free unless notFree names it), as the router asks them. That a machine is
// root-free is the controller's own judgement, proven in its repository (probe_agent_root_test.go); here the
// lab says it, so the proof can show what the router does with each answer.
func (l *lab) playController() {
l.t.Helper()
free, err := l.observer.Subscribe("mesh.seat.mesh-controller.tool.root-free", func(m *nats.Msg) {
var args struct {
Machines []string `json:"machines"`
}
_ = json.Unmarshal(m.Data, &args)
var out []map[string]any
l.mu.Lock()
for _, name := range args.Machines {
why, not := l.notFree[name]
if !not {
why = "the lab says no agent can become root here"
}
out = append(out, map[string]any{"machine": name, "free": !not, "why": why,
"judged": time.Now().UTC().Format(time.RFC3339Nano)})
}
l.rootAsked++
l.mu.Unlock()
raw, _ := json.Marshal(map[string]any{"result": map[string]any{"machines": out}, "node": machine})
_ = m.Respond(raw)
})
if err != nil {
l.t.Fatal(err)
}
l.t.Cleanup(func() { _ = free.Unsubscribe() })
sub, err := l.observer.Subscribe("mesh.seat.mesh-controller.tool.conditions", func(m *nats.Msg) {
l.mu.Lock()
list := append([]map[string]any{}, l.conditions...)
l.mu.Unlock()
raw, _ := json.Marshal(map[string]any{"result": map[string]any{"conditions": list}, "node": machine})
_ = m.Respond(raw)
})
if err != nil {
l.t.Fatal(err)
}
l.t.Cleanup(func() { _ = sub.Unsubscribe() })
}
// raise says a condition as the controller does: on its event, and in its conditions verb from then on.
func (l *lab) raise(c map[string]any) {
l.t.Helper()
l.mu.Lock()
l.conditions = append(l.conditions, c)
l.mu.Unlock()
l.sayCondition("condition-raised", c)
}
// clear ends a condition as the controller does.
func (l *lab) clear(key string) {
l.t.Helper()
l.mu.Lock()
var kept []map[string]any
var gone map[string]any
for _, c := range l.conditions {
if c["key"] == key {
gone = c
continue
}
kept = append(kept, c)
}
l.conditions = kept
l.mu.Unlock()
if gone != nil {
gone["cleared"] = time.Now().UTC().Format(time.RFC3339Nano)
l.sayCondition("condition-cleared", gone)
}
}
func (l *lab) sayCondition(event string, c map[string]any) {
l.t.Helper()
nc, _ := l.as("controller")
js, err := nc.JetStream()
if err != nil {
l.t.Fatal(err)
}
raw, _ := json.Marshal(c)
msg := nats.NewMsg("mesh.seat.mesh-controller.event." + event)
msg.Data = raw
msg.Header.Set("x-event-id", fmt.Sprintf("%s.%v.%d", event, c["key"], time.Now().UnixNano()))
if _, err := js.PublishMsg(msg); err != nil {
l.t.Fatalf("saying %s as the controller: %v", event, err)
}
}
// playDesk is the desk channel (kind `desktop`) at the bus, on the machine's runtime's credential, which
// carries the desk channel: it says it is ready, and keeps what the router shows it — the link's code.
func (l *lab) playDesk(ctx context.Context) {
l.t.Helper()
// The desk channel runs as an account of its own (it carries a link's code), on its own credential.
nc, _ := l.as(machine + ".desk-channel")
stand := func() {
raw, _ := json.Marshal(map[string]any{"ready": true, "at": time.Now().UTC()})
js, _ := nc.JetStream()
msg := nats.NewMsg("mesh.seat.intake.event.standing.desktop")
msg.Data = raw
msg.Header.Set("x-event-id", fmt.Sprintf("standing.desktop.%d", time.Now().UnixNano()))
if _, err := js.PublishMsg(msg); err != nil {
l.t.Logf("the desk could not say it is ready: %v", err)
}
}
stand()
go func() {
tick := time.NewTicker(20 * time.Second)
defer tick.Stop()
for {
select {
case <-ctx.Done():
return
case <-tick.C:
stand()
}
}
}()
js, err := jetstream.New(nc)
if err != nil {
l.t.Fatal(err)
}
bind, cancel := context.WithTimeout(ctx, 10*time.Second)
defer cancel()
worker, err := js.Consumer(bind, "SEAT_CHANNEL", "SEAT_CHANNEL_DESKTOP_worker")
if err != nil {
l.t.Fatalf("the desk's worker: %v", err)
}
consuming, err := worker.Consume(func(m jetstream.Msg) {
var shown struct {
Title string `json:"title"`
Body string `json:"body"`
}
_ = json.Unmarshal(m.Data(), &shown)
l.mu.Lock()
l.desk = append(l.desk, shown.Title+": "+shown.Body)
l.mu.Unlock()
_ = m.Ack()
})
if err != nil {
l.t.Fatal(err)
}
l.t.Cleanup(consuming.Stop)
}
var linkCode = regexp.MustCompile(`(\d{3}) (\d{3})\.`)
// codeOnTheDesk waits for the link's code the router shows on the desk.
func (l *lab) codeOnTheDesk(within time.Duration) string {
l.t.Helper()
deadline := time.Now().Add(within)
for time.Now().Before(deadline) {
l.mu.Lock()
for _, s := range l.desk {
if m := linkCode.FindStringSubmatch(s); m != nil {
l.mu.Unlock()
return m[1] + m[2]
}
}
l.mu.Unlock()
time.Sleep(100 * time.Millisecond)
}
l.mu.Lock()
defer l.mu.Unlock()
l.t.Fatalf("the desk was never shown a code; it was shown %q\nthe router said:\n%s", l.desk, l.logs["messenger"])
return ""
}
// routerRecord reads one key of the router's state, as the lab's observer.
func (l *lab) routerRecord(bucket, key string) (jetstream.KeyValueEntry, jetstream.KeyValue) {
l.t.Helper()
js, err := jetstream.New(l.observer)
if err != nil {
l.t.Fatal(err)
}
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
kv, err := js.KeyValue(ctx, bucket)
if err != nil {
l.t.Fatalf("the router's bucket %s: %v", bucket, err)
}
e, err := kv.Get(ctx, key)
if err != nil {
l.t.Fatalf("the router's %s %s: %v", bucket, key, err)
}
return e, kv
}
+68
View File
@@ -0,0 +1,68 @@
#!/bin/sh
# The live acceptance of the operator's answers (novox/hq ADR 0259), after rollout, at the controller's
# terminal on the control node. The lab's proof (TestTheOperatorsAnswerEndToEnd) ran the same flow with a fake
# Telegram; this runs it once for real, with a question that changes nothing.
#
# sh live-acceptance.sh check, ask the rehearsal, wait for the answer, read the record
#
# It reads and asks only: no setting, assignment or push. The one thing it starts is a rehearsal, a question the
# operator answers on the phone; approving it performs nothing and is recorded as the operator's decision.
#
# It stops at the first check that fails and says what to do. Exit 0 means: the question reached the phone
# with Approve and Decline, the answer came back as a warrant, and the hand-act log names who answered,
# through which channel, with which proofs.
set -eu
# mesh-cli asks the controller; on the control node, from the operator's own account, a line is the
# controller's terminal (ADR 0272). `nox` is its interactive alias; a script says mesh-cli.
mc=${MESH_CLI:-mesh-cli}
stop() { echo "NOT ACCEPTED: $*"; exit 1; }
echo "1. Nothing open says the router's or Telegram's machine is not root-free (root-not-free, agent-root)."
if $mc conditions | grep -q 'root-not-free\|agent-root'; then
$mc conditions | grep 'root-not-free\|agent-root'
stop "while an agent can become root there, Telegram only acknowledges; close it first (ADR 0266, ADR 0268)"
fi
echo "2. The Telegram channel holds a working bot token and its updates arrive."
tg=$($mc ask telegram telegram_status) || stop "the telegram module does not answer"
echo "$tg" | grep -q '"can_send": *true' || { echo "$tg"; stop "Telegram cannot send now (why_not above): give the bot token again at the controller's terminal (hq ADR 0259, the operator's steps), then push the control node"; }
echo "3. One Telegram account is linked as the operator, and its hour has passed."
ids=$($mc ask messenger messenger_identities) || stop "the router does not answer"
echo "$ids"
echo "$ids" | grep -q '"telegram"' || stop "no Telegram account is linked: send /start to the bot, then the code from the desk"
# The router answers approvals from a new link only an hour after it was made (ADR 0259 §7): said here, with the
# time, rather than seen as a rehearsal nobody's approval reaches.
from=$(echo "$ids" | tr -d '\n ' | grep -o '"kind":"telegram"[^}]*' | grep -o '"approve-from":"[^"]*"' | cut -d'"' -f4 | head -n 1)
[ -n "$from" ] || stop "the router's word on the Telegram account says no time it can approve from"
starts=$(date -d "$from" +%s 2>/dev/null) || stop "the time the Telegram account can approve from is not a time: $from"
if [ "$(date +%s)" -lt "$starts" ]; then
stop "the Telegram account was linked less than an hour ago — approvals start at $(date -d "$from" '+%H:%M') (local time); run this again then"
fi
echo "4. The rehearsal: a question on your phone. Approve it there within 15 minutes."
said=$($mc rehearse --for 15m) || stop "the rehearsal could not be asked"
echo "$said"
id=$(echo "$said" | sed -n 's/^rehearsal \([A-Za-z0-9_-]*\) asked.*/\1/p')
[ -n "$id" ] || stop "the rehearsal did not say its id"
echo "5. Waiting for your answer in the hand-act log (every 10 s, at most 16 minutes)."
i=0
while [ $i -lt 96 ]; do
acts=$($mc hand-acts --days 1 --json 2>/dev/null || true)
mine=$(echo "$acts" | tr -d '\n' | grep -o "{[^{}]*\"ask\": *\"$id\"[^{}]*}" || true)
if [ -n "$mine" ]; then
echo "$mine"
echo "$mine" | grep -q 'rehearsal=decline' && stop "the rehearsal was declined: run it again and choose Approve to see an approval recorded"
echo "$mine" | grep -q 'rehearsal=approve' || stop "the rehearsal's record names no approval"
echo "$mine" | grep -q '"via": *"telegram (telegram), user id verified"' || stop "the record does not say the answer came through Telegram from a verified account"
echo "$mine" | grep -q '"P1"' || stop "the record does not carry the proof"
echo "$mine" | grep -q '"outcome": *"done"' || stop "the rehearsal's act did not end done"
echo "ACCEPTED: the rehearsal was approved on Telegram, acted on once, and recorded as your decision."
exit 0
fi
i=$((i + 1))
sleep 10
done
stop "no answer was recorded within 16 minutes: check messenger_status and telegram_status"
+37
View File
@@ -0,0 +1,37 @@
#!/bin/sh
# mesh-check-toolchain: go
#
# The proof of the operator's answers (novox/hq ADR 0259), the Go part of the lab's own check: run by the
# build seat in the mesh's Go toolchain, its own container, because the lab's merge-check.sh declares it
# (`# mesh-check-also: go asks/merge-check.sh`) — and by hand: `sh asks/merge-check.sh`.
#
# Formatted, vetted, and run against the clones beside the lab: on the build seat those are the controller,
# the runtime and the catalogue **the mesh runs**, so every pull request of the lab proves the answers the
# operator would get today. By hand, MESH_LAB_REPOS, or MESH_LAB_CONTROLLER, MESH_LAB_TOOLS and
# MESH_LAB_CATALOGUE, point it at branches under review. It takes about four minutes when it runs.
#
# **Said, never passed silently**: where a clone is missing, or the controller predates the proof's fixture,
# it says NOT RUN and why, in capitals, and passes; a proof that runs and fails, fails the check.
set -eu
cd "$(dirname "$0")"
export GOPRIVATE=git.novox.be
unformatted=$(gofmt -l .)
if [ -n "$unformatted" ]; then
echo "not gofmt'd:"
echo "$unformatted"
exit 1
fi
sh -n live-acceptance.sh
go vet ./...
out=$(mktemp)
status=0
go test -count=1 -timeout 25m -v -run '^TestTheOperatorsAnswerEndToEnd$' . >"$out" 2>&1 || status=$?
grep -E '^(--- | --- |ok|FAIL|PASS)' "$out" || true
if grep -q 'NOT RUN' "$out"; then
echo "NOT RUN: THE PROOF OF THE OPERATOR'S ANSWERS DID NOT RUN HERE:"
grep -o 'NOT RUN: .*' "$out" | head -3
fi
[ "$status" -eq 0 ] || tail -60 "$out"
rm -f "$out"
exit "$status"
+347
View File
@@ -0,0 +1,347 @@
package asks
import (
"context"
"encoding/json"
"fmt"
"os"
"path/filepath"
"strings"
"testing"
"time"
"git.novox.be/novox/mesh-sdk/go/asker"
"git.novox.be/novox/mesh-sdk/go/asks"
)
// The operator's own Telegram account in the lab, and somebody else's.
const (
operatorAccount = 42
strangerAccount = 99
groupChat = -1001
)
// TestTheOperatorsAnswerEndToEnd is the operator's acceptance of novox/hq ADR 0259, run before anybody uses it:
//
// "I get a Telegram message from the mesh with what is asked, why, by whom, and Approve/Decline buttons;
// tapping Approve makes the mesh do exactly that one thing once; Decline or no answer does nothing;
// nobody else (no agent, no other Telegram user, no replay of an old message) can produce an approval;
// and I can see in the record who approved what."
//
// What is real: the bus, of the release the mesh runs, with the accounts the controller at its commit composes
// for one machine and raised by that controller's own code (its fixture TestTheAsksLabBus); the mesh's runtime
// (mesh-tools node-tools) serving the router (messenger) and the Telegram channel (telegram), each on its own
// credential, as `runs-as` places them; the router's and the channel's code at their commits; the SDK's asker.
//
// What is the lab's: Telegram itself (a fake Bot API, the operator's phone), the desk (the desktop kind at the
// bus, which shows the link's code), the controller's conditions verb (a list the lab keeps), and the asker
// (the module `lab-asker`, the SDK's client in this process). The controller as an asker — that it acts on a
// warrant once, through the verb the action names, and records it — is its own repository's test
// (asker_test.go); the live acceptance after rollout runs it for real.
func TestTheOperatorsAnswerEndToEnd(t *testing.T) {
repos := clonesFor(t)
dir := os.Getenv("MESH_LAB_ASKS_KEEP")
if dir == "" {
dir = t.TempDir()
} else if err := os.MkdirAll(dir, 0o700); err != nil {
t.Fatal(err)
}
l := &lab{t: t, repos: repos, dir: dir}
token := "123456789:" + strings.Repeat("A", 35)
l.tg = newFakeTelegram(t, token)
// The bus, as the controller composes and raises it.
l.compose("")
l.startBus()
l.compose(l.url)
l.playController()
// The programs, at their commits. The channel is built to call the lab's phone in place of Telegram.
runtimeBin := l.build(filepath.Join(repos.tools, "node-tools"), "./cmd/node-tools", "node-tools", "")
routerBin := l.build(filepath.Join(repos.catalogue, "modules", "messenger"), "./cmd/messenger", "messenger", "")
channelBin := l.build(filepath.Join(repos.catalogue, "modules", "telegram"), "./cmd/telegram", "telegram",
"-X main.API="+l.tg.URL())
routerState := filepath.Join(dir, "messenger")
if err := os.MkdirAll(routerState, 0o700); err != nil {
t.Fatal(err)
}
if err := os.WriteFile(filepath.Join(routerState, "settings.json"), []byte(`{"time-zone": "UTC"}`), 0o600); err != nil {
t.Fatal(err)
}
tokenFile := filepath.Join(dir, "telegram", "telegram-token")
if err := os.MkdirAll(filepath.Dir(tokenFile), 0o700); err != nil {
t.Fatal(err)
}
if err := os.WriteFile(tokenFile, []byte(token+"\n"), 0o600); err != nil {
t.Fatal(err)
}
ctx, stop := context.WithCancel(context.Background())
t.Cleanup(stop)
l.playDesk(ctx)
l.runtime("messenger", runtimeBin, routerBin, map[string]string{
"MESH_MESSENGER_SETTINGS": filepath.Join(routerState, "settings.json"), "MESH_MESSENGER_STATE": routerState})
l.runtime("telegram", runtimeBin, channelBin, map[string]string{"MESH_TELEGRAM_TOKEN_FILE": tokenFile})
l.waitSaid("messenger", "taking the asks submitted", time.Minute)
l.waitSaid("telegram", "taking the router's messages for this kind", time.Minute)
a := newLabAsker(t, l)
// --- 1. The operator links their account: /start on the phone, the code from the desk.
t.Run("the operator links their account with the code shown on the desk", func(t *testing.T) {
since := time.Now()
l.tg.typed(operatorAccount, operatorAccount, "private", "/start")
code := l.codeOnTheDesk(90 * time.Second)
l.tg.typed(operatorAccount, operatorAccount, "private", code[:3]+" "+code[3:])
l.tg.waitFor("the link said done", operatorAccount, since, 90*time.Second, func(m tgMessage) bool {
return strings.Contains(m.Text, "Linked")
})
})
if t.Failed() {
return
}
// --- 2. Within the hour after a link, an approval is refused; the button is not spent.
var first asks.Ask
var firstMessage tgMessage
t.Run("an approval in the hour after linking is refused and nothing is done", func(t *testing.T) {
since := time.Now()
first = a.ask("a1", "Flip the lab's switch?", time.Hour)
firstMessage = l.tg.waitFor("the first question with its buttons", operatorAccount, since, 90*time.Second,
func(m tgMessage) bool { return !m.Edited && button(m, "Approve") != "" && button(m, "Decline") != "" })
for _, says := range []string{"Flip the lab's switch?", "approve or decline", "lab-asker",
"the switch is flipped, once", "nothing is flipped"} {
if !strings.Contains(firstMessage.Text, says) {
t.Errorf("the question does not say %q (what is asked, why, by whom, what each answer does, what "+
"no answer does): %q", says, firstMessage.Text)
}
}
tapped := time.Now()
l.tg.tapped(operatorAccount, operatorAccount, "private", firstMessage.ID, button(firstMessage, "Approve"))
l.tg.waitFor("the refusal of an approval inside the hour", operatorAccount, tapped, 90*time.Second,
func(m tgMessage) bool {
return strings.Contains(m.Text, "That answer was not taken") && strings.Contains(m.Text, "linked at")
})
if !a.noWordOn("a1", 2*time.Second) || len(a.performed()) != 0 {
t.Fatalf("an approval inside the hour was a warrant: %v", a.performed())
}
})
// The hour passes: the lab ages the link in the router's own record (the hour itself is the router's
// unit test: an approve at 59 minutes refused, at 60 taken).
t.Run("the hour after the link passes", func(t *testing.T) {
e, kv := l.routerRecord("messenger_identities", fmt.Sprintf("telegram.%d", operatorAccount))
var id map[string]any
if err := json.Unmarshal(e.Value(), &id); err != nil {
t.Fatal(err)
}
past := time.Now().Add(-61 * time.Minute).UTC()
id["linked"], id["approve-from"] = past.Format(time.RFC3339Nano), past.Add(time.Hour).Format(time.RFC3339Nano)
raw, _ := json.Marshal(id)
if _, err := kv.Update(context.Background(), e.Key(), raw, e.Revision()); err != nil {
t.Fatal(err)
}
})
// --- 3. Approve: the warrant, for this ask alone, performed once.
t.Run("approve does exactly the one thing, once, and says who approved", func(t *testing.T) {
tapped := time.Now()
l.tg.tapped(operatorAccount, operatorAccount, "private", firstMessage.ID, button(firstMessage, "Approve"))
w := a.warrantFor("a1", 30*time.Second)
if w.Outcome != asks.OutcomeChosen || w.Option != "approve" || w.Level != asks.Approve {
t.Fatalf("the warrant: %+v", w)
}
if w.By == nil || w.By.Kind != "telegram" || w.By.Identity != fmt.Sprint(operatorAccount) ||
w.By.Verified != "user id verified" || w.Channel != "telegram" || strings.Join(w.Proofs, ",") != "P1" {
t.Errorf("the warrant does not say who approved, how and where: %+v by %+v", w, w.By)
}
if w.AskDigest != first.Digest() {
t.Errorf("the warrant names the digest %s, not the ask's %s", w.AskDigest, first.Digest())
}
if got := a.performed(); len(got) != 1 || got[0]["arg.switch"] != "approve" {
t.Fatalf("performed %v, want the one bound act once", got)
}
// Every copy says how it ended, its buttons gone.
l.tg.waitFor("the question edited to its outcome", operatorAccount, tapped, 90*time.Second,
func(m tgMessage) bool { return m.Edited && m.ID == firstMessage.ID && len(m.Buttons) == 0 })
// And the router's record holds who approved what.
e, _ := l.routerRecord("messenger_asks", "lab-asker.a1")
var rec struct {
State string `json:"state"`
Warrant *asks.Warrant `json:"warrant"`
}
_ = json.Unmarshal(e.Value(), &rec)
if rec.State != "chosen" || rec.Warrant == nil || rec.Warrant.By == nil || rec.Warrant.By.Identity != fmt.Sprint(operatorAccount) {
t.Errorf("the router's record: %s", e.Value())
}
t.Logf("the record reads: %s", w.Says())
})
// --- 4. A replayed tap, and a redelivered warrant, do nothing more.
t.Run("a replayed tap and a redelivered warrant do nothing more", func(t *testing.T) {
tapped := time.Now()
l.tg.tapped(operatorAccount, operatorAccount, "private", firstMessage.ID, button(firstMessage, "Approve"))
l.tg.waitFor("the refusal of a replayed tap", operatorAccount, tapped, 90*time.Second,
func(m tgMessage) bool { return strings.Contains(m.Text, "That answer was not taken") })
time.Sleep(2 * time.Second)
if n := a.wordsOn("a1"); n != 1 {
t.Errorf("the router said %d words on one ask", n)
}
w := a.warrantFor("a1", time.Second)
if h := a.take(w); h != asker.Done {
t.Errorf("a warrant heard again was %s", h)
}
if got := a.performed(); len(got) != 1 {
t.Fatalf("performed %v after a replay", got)
}
})
// --- 5. Somebody else's account, and a group, cannot answer; Decline does nothing.
t.Run("another account and a group cannot answer, and decline does nothing", func(t *testing.T) {
since := time.Now()
a.ask("a2", "Flip it once more?", time.Hour)
m := l.tg.waitFor("the second question", operatorAccount, since, 90*time.Second,
func(m tgMessage) bool { return !m.Edited && button(m, "Approve") != "" })
// The buttons' data reached a stranger (a forwarded question), who taps in their own chat with the bot:
// the channel never showed that message there, so the tap goes nowhere, and the stranger is told
// nothing. (A stranger's answer that does reach the router is refused there and told to the operator by
// name: the router's own test.)
l.tg.tapped(strangerAccount, strangerAccount, "private", m.ID, button(m, "Approve"))
// The operator taps a message whose words are not the ones the channel put there.
l.tg.tappedOn(operatorAccount, operatorAccount, "private", m.ID, button(m, "Approve"), "Flip ALL the switches?")
// The operator taps in a group the bot is in: not their own chat, dropped by the channel.
l.tg.tapped(operatorAccount, groupChat, "supergroup", m.ID, button(m, "Approve"))
if !a.noWordOn("a2", 3*time.Second) {
t.Fatal("a tap from another account, on changed words, or from a group, was a warrant")
}
if got := l.tg.since(strangerAccount, since); len(got) != 0 {
t.Errorf("the bot answered the stranger: %+v", got)
}
if got := l.tg.since(groupChat, since); len(got) != 0 {
t.Errorf("the bot answered in a group: %+v", got)
}
l.tg.tapped(operatorAccount, operatorAccount, "private", m.ID, button(m, "Decline"))
w := a.warrantFor("a2", 30*time.Second)
if w.Option != "decline" || w.Level != asks.Acknowledge {
t.Errorf("decline's warrant: %+v", w)
}
if got := a.performed(); len(got) != 1 {
t.Fatalf("decline performed something: %v", got)
}
})
// --- 6. No agent can produce an approval: the bus refuses every forgery, as the controller composed it.
t.Run("no agent, channel or module but the router can say a warrant, and none asks in another's name", func(t *testing.T) {
warrant, _ := json.Marshal(asks.Warrant{Ask: "a9", Asker: "lab-asker", Outcome: asks.OutcomeChosen,
Option: "approve", Level: asks.Approve, AskDigest: first.Digest(),
By: &asks.Person{Who: asks.Operator, Kind: "telegram", Identity: "42", Verified: "user id verified"}})
decided := asks.DecidedSubject("lab-asker")
// The control: what is granted is not refused, so a refusal below is the bus's word and not the lab's.
// The channel's control is a true standing, now: a false one would hold the router's messages.
standing, _ := json.Marshal(asks.Standing{Ready: true, EditsSilently: true, At: time.Now().UTC()})
for _, try := range []struct {
user, subject string
body []byte
}{
{machine + ".lab-asker", "mesh.seat.operator-channel.accept.cancel.lab-asker", []byte(`{"id":"none"}`)},
{machine + ".telegram", asks.IntakeSubject(asks.StandingWhat, "telegram"), standing},
} {
if l.refusedToPublish(try.user, try.subject, try.body) {
t.Fatalf("the control failed: %s was refused %s, which the controller grants it", try.user, try.subject)
}
}
for _, user := range []string{
machine + ".node-tools", // the machine's runtime: every agent's tool calls go through one
machine + ".lab-asker", // the asker itself
machine + ".lab-bystander", // any other module
machine + ".telegram", // the channel
"node." + machine, // the machine's node-engine
"controller", // the controller asks; it does not decide
} {
if !l.refusedToPublish(user, decided, warrant) {
t.Errorf("%s may say a warrant on %s", user, decided)
}
}
for _, try := range []struct{ user, subject, what string }{
{machine + ".lab-asker", asks.AskSubject("mesh-controller"), "an ask in another asker's name"},
{machine + ".lab-bystander", asks.AskSubject("lab-bystander"), "an ask from a module that may not ask"},
{machine + ".node-tools", asks.AskSubject("mesh-controller"), "an agent asking as the controller through the runtime"},
{machine + ".node-tools", asks.IntakeSubject(asks.Chosen, "telegram"), "an agent saying a Telegram tap"},
{machine + ".lab-asker", asks.IntakeSubject(asks.Chosen, "telegram"), "an asker saying a Telegram tap"},
{machine + ".node-tools", "$KV.messenger_identities.telegram.666", "an agent linking an account of its own"},
{machine + ".lab-asker", "$KV.messenger_asks.lab-asker.a1", "an asker rewriting the router's record"},
{machine + ".telegram", asks.ProofSubject(asks.CodeProof, "desktop"), "a channel speaking for another kind"},
} {
if !l.refusedToPublish(try.user, try.subject, []byte(`{}`)) {
t.Errorf("%s: %s may publish %s", try.what, try.user, try.subject)
}
}
time.Sleep(time.Second)
if n := a.wordsOn("a9"); n != 0 || len(a.performed()) != 1 {
t.Fatalf("a forged warrant reached the asker: %d words, performed %v", n, a.performed())
}
})
// --- 7. While an agent can become root where the router runs, nothing proven there approves.
t.Run("while the controller cannot say root is free where the router runs, nothing approves", func(t *testing.T) {
// An approval the operator already sees on the phone, then root stops being free.
since := time.Now()
a.ask("a5", "Flip it while root is checked?", time.Hour)
m := l.tg.waitFor("the question asked while root was free", operatorAccount, since, 90*time.Second,
func(m tgMessage) bool { return !m.Edited && button(m, "Approve") != "" })
l.mu.Lock()
l.notFree = map[string]string{machine: "an agent can become root without a person"}
asked := l.rootAsked
l.mu.Unlock()
tapped := time.Now()
l.tg.tapped(operatorAccount, operatorAccount, "private", m.ID, button(m, "Approve"))
l.tg.waitFor("the refusal of an approval while root is not free", operatorAccount, tapped, 90*time.Second,
func(m tgMessage) bool { return strings.Contains(m.Text, "That answer was not taken") })
l.mu.Lock()
askedAgain := l.rootAsked > asked
l.mu.Unlock()
if !askedAgain {
t.Error("the router judged an approval without asking the controller whether root is free")
}
if !a.noWordOn("a5", 2*time.Second) {
t.Fatal("an approval was a warrant while root was not free")
}
// A new ask that needs an approval is refused to its asker, in words: no channel can carry it now.
a.ask("a3", "Flip it while root is not free?", time.Hour)
w := a.warrantFor("a3", 30*time.Second)
if w.Outcome != asks.OutcomeRefused || w.By != nil {
t.Errorf("an approving ask while root is not free: %+v", w)
}
if got := a.performed(); len(got) != 1 {
t.Fatalf("performed %v while root was not free", got)
}
l.mu.Lock()
l.notFree = nil
l.mu.Unlock()
if err := a.client.Cancel("a5"); err != nil {
t.Errorf("taking the ask back: %v", err)
}
})
// --- 8. No answer does nothing: the ask expires.
t.Run("an ask nobody answers expires and nothing is done", func(t *testing.T) {
since := time.Now()
a.ask("a4", "Flip it if you answer in time?", 20*time.Second)
m := l.tg.waitFor("the fourth question", operatorAccount, since, 90*time.Second,
func(m tgMessage) bool { return !m.Edited && button(m, "Approve") != "" })
w := a.warrantFor("a4", 2*time.Minute)
if w.Outcome != asks.OutcomeExpired || w.By != nil {
t.Errorf("the end of an unanswered ask: %+v", w)
}
late := time.Now()
l.tg.tapped(operatorAccount, operatorAccount, "private", m.ID, button(m, "Approve"))
l.tg.waitFor("the refusal of a late tap", operatorAccount, late, 90*time.Second,
func(m tgMessage) bool { return strings.Contains(m.Text, "That answer was not taken") })
time.Sleep(2 * time.Second)
if got := a.performed(); len(got) != 1 {
t.Fatalf("an expired ask performed something: %v", got)
}
})
t.Logf("performed, over the whole run: %v (one approval, one act)", a.performed())
}
+7 -3
View File
@@ -1,6 +1,7 @@
#!/bin/sh
# mesh-check-toolchain: typescript
# mesh-check-also: go replays/merge-check.sh
# mesh-check-also: go asks/merge-check.sh
#
# The lab's own check (novox/hq ADR 0237 as amended): the second layer of a pull request's merge check,
# `mesh/repo-check`, run by the build seat — and by hand.
@@ -11,8 +12,10 @@
#
# - this script, in the TypeScript toolchain: installed from the lock file, type-checked, sources and
# tests, and the unit suite;
# - replays/merge-check.sh, in the Go toolchain, declared on the line above: the replays register and its
# prover formatted, vetted and compiled, and the register's own tests.
# - replays/merge-check.sh, in the Go toolchain, declared above: the replays register and its prover
# formatted, vetted and compiled, and the register's own tests;
# - asks/merge-check.sh, in the Go toolchain, declared above: the proof of the operator's answers, run
# against the controller, runtime and catalogue cloned beside the lab (novox/hq ADR 0259).
#
# **Said, never passed silently**: the integration suite raises a lab on a workstation, and the replays
# raise containers through the container runtime — neither is given to code nobody has approved. The
@@ -27,7 +30,8 @@ npm test
# By hand, with Go at hand, the Go part runs here too; on the build seat it runs in the Go toolchain.
if command -v go >/dev/null 2>&1; then
sh replays/merge-check.sh
sh asks/merge-check.sh
else
echo "replays/ (Go) is checked by replays/merge-check.sh, which the build seat runs in the Go toolchain"
echo "replays/ and asks/ (Go) are checked by their merge-check.sh, which the build seat runs in the Go toolchain"
fi
echo "NOT RUN HERE: the integration suite and the replays — they need a lab and the container runtime"