Merge pull request 'Prove the operator's answers end to end before anybody uses them (hq ADR 0259)' (#74) from proofs/asks-answered-on-the-phone into main
This commit was merged in pull request #74.
This commit is contained in:
@@ -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
|
||||
}
|
||||
@@ -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
|
||||
@@ -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
@@ -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
@@ -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=
|
||||
@@ -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
|
||||
}
|
||||
@@ -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"
|
||||
@@ -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"
|
||||
@@ -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
@@ -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"
|
||||
|
||||
Reference in New Issue
Block a user