mesh/merge-gate pass: the change touches no module of the mesh's graph
mesh/repo-check pass: THE CHANGE ALTERS ITS OWN CHECK (merge-check.sh): main's version judged it; the change's judges the pull requests after it merges; it…
mesh/delivery superseded: a newer head of the same pull request
538 lines
17 KiB
Go
538 lines
17 KiB
Go
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
|
|
}
|