Merge pull request 'Phase 5: replay the core incidents, and prove each fails before its fix (hq ADR 0237)' (#52) from feat/replays into main
mesh/delivery held for a person: merged without a passing check: only a person decides that it goes on
mesh/delivery held for a person: merged without a passing check: only a person decides that it goes on
This commit was merged in pull request #52.
This commit is contained in:
@@ -406,3 +406,20 @@ That suite earned itself on its first run: it found that a snapshot of a running
|
||||
miss a file written seconds earlier — not stale, **absent** — because the write was still in
|
||||
the guest's page cache. The design had listed that as an open question. The test answered it,
|
||||
and `snapshot` now flushes first.
|
||||
|
||||
## The replays (`replays/`)
|
||||
|
||||
Every core incident, replayed (novox/hq to-be 45 §9, M9) — in Go, a module of its own. `register.go`
|
||||
names each replay, its issue, its fix and how it is run; a replay of one component's logic is a test in
|
||||
that component's repository (the controller's `replays_test.go`), and a replay of what the mesh runs —
|
||||
the bus server's release, the resolver the catalogue configures under musl and glibc — is here.
|
||||
|
||||
```
|
||||
cd replays
|
||||
MESH_TEST_NATS=nats://… go test ./... the bus and resolver replays (the resolver's needs a container runtime)
|
||||
go run ./cmd/prove [R262 …] each replay on the commit before its fix (must fail) and on it (must pass)
|
||||
```
|
||||
|
||||
The prover finds the core repositories beside this one, as the beds do. The build seat runs the bus and
|
||||
resolver replays before every core and catalogue merge, against the bus of the release the mesh runs and
|
||||
the change's own catalogue.
|
||||
|
||||
@@ -0,0 +1,141 @@
|
||||
package replays
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"math/rand"
|
||||
"os"
|
||||
"sync"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/nats-io/nats.go"
|
||||
)
|
||||
|
||||
// **R266 — a consumer with several filters is handed every message** (novox/hq issue 266).
|
||||
//
|
||||
// The controller follows the forge's merges, build outcomes and providers' standings on one durable
|
||||
// consumer with several filters. On nats 2.10.29 such a consumer was moved past a message now and then
|
||||
// without handing it over — nothing pending, nothing redelivered — and on 2026-10-06 that message was a
|
||||
// merge nobody acted on. This is that consumer under traffic shaped like the mesh's, against the bus
|
||||
// MESH_TEST_NATS names: in a merge check, a throwaway of the release the mesh runs; for the prover, the
|
||||
// release the catalogue pinned at the commit it proves. Fails on 2.10.29, passes on 2.11.17.
|
||||
//
|
||||
// Everything it makes is under a prefix of its own, so it can run on a bus other things use.
|
||||
func TestReplay266AConsumerWithSeveralFiltersIsHandedEveryMessage(t *testing.T) {
|
||||
url := os.Getenv("MESH_TEST_NATS")
|
||||
if url == "" {
|
||||
t.Skip("MESH_TEST_NATS names no bus to replay against")
|
||||
}
|
||||
p := fmt.Sprintf("r266x%d", time.Now().UnixNano())
|
||||
admin, err := nats.Connect(url)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer admin.Close()
|
||||
js, err := admin.JetStream()
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
stream := "R266_" + p
|
||||
if _, err := js.AddStream(&nats.StreamConfig{Name: stream, Subjects: []string{p + ".mod.*.event.>", p + ".seat.*.event.>"},
|
||||
MaxAge: time.Hour, MaxMsgsPerSubject: 10000, Storage: nats.FileStorage}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer func() { _ = js.DeleteStream(stream) }()
|
||||
// The controller's consumer, as the mesh declared it on the day: seven filters, one at a time.
|
||||
if _, err := js.AddConsumer(stream, &nats.ConsumerConfig{Durable: "controller", DeliverSubject: "_DELIVER." + p,
|
||||
FilterSubjects: []string{
|
||||
p + ".mod.mesh-catalog.event.upgraded",
|
||||
p + ".mod.mesh-catalog.event.catching-up",
|
||||
p + ".seat.node-build-agent.event.built",
|
||||
p + ".mod.gitea.event.pull.merged",
|
||||
p + ".seat.mesh-build-machine.event.built",
|
||||
p + ".mod.*.event.provisioner.failing",
|
||||
p + ".mod.*.event.provisioner.recovered",
|
||||
},
|
||||
AckPolicy: nats.AckExplicitPolicy, AckWait: 30 * time.Second, MaxDeliver: 5, MaxAckPending: 1,
|
||||
DeliverPolicy: nats.DeliverNewPolicy}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
controller, err := nats.Connect(url)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer controller.Close()
|
||||
cjs, _ := controller.JetStream()
|
||||
handed := make(chan *nats.Msg, 64)
|
||||
sub, err := cjs.ChanSubscribe("", handed, nats.Bind(stream, "controller"))
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer func() { _ = sub.Unsubscribe() }()
|
||||
var got sync.Map
|
||||
go func() {
|
||||
for m := range handed {
|
||||
if meta, err := m.Metadata(); err == nil {
|
||||
got.Store(meta.Sequence.Stream, true)
|
||||
}
|
||||
time.Sleep(time.Duration(rand.Intn(20)) * time.Millisecond)
|
||||
_ = m.Ack()
|
||||
}
|
||||
}()
|
||||
|
||||
followed := []string{p + ".mod.gitea.event.pull.merged", p + ".seat.node-build-agent.event.built",
|
||||
p + ".mod.postgres.event.provisioner.recovered", p + ".mod.mesh-catalog.event.catching-up"}
|
||||
var mu sync.Mutex
|
||||
var sent []uint64
|
||||
stop := time.Now().Add(6 * time.Second)
|
||||
var wg sync.WaitGroup
|
||||
for w := 0; w < 4; w++ {
|
||||
wg.Add(1)
|
||||
go func(w int) {
|
||||
defer wg.Done()
|
||||
nc, err := nats.Connect(url)
|
||||
if err != nil {
|
||||
t.Error(err)
|
||||
return
|
||||
}
|
||||
defer nc.Close()
|
||||
pjs, _ := nc.JetStream()
|
||||
for i := 0; time.Now().Before(stop); i++ {
|
||||
if rand.Intn(40) == 0 {
|
||||
if ack, err := pjs.Publish(followed[rand.Intn(len(followed))], []byte(`{}`)); err == nil {
|
||||
mu.Lock()
|
||||
sent = append(sent, ack.Sequence)
|
||||
mu.Unlock()
|
||||
}
|
||||
} else {
|
||||
// A build's log, on a subject of its own: what the stream mostly holds.
|
||||
_, _ = pjs.Publish(fmt.Sprintf("%s.seat.node-build-agent.event.log.build-%d-%d", p, w, i/50), make([]byte, 200))
|
||||
}
|
||||
if rand.Intn(100) == 0 {
|
||||
time.Sleep(time.Duration(rand.Intn(300)) * time.Millisecond)
|
||||
}
|
||||
}
|
||||
}(w)
|
||||
}
|
||||
wg.Wait()
|
||||
if len(sent) == 0 {
|
||||
t.Fatal("no followed event was published")
|
||||
}
|
||||
|
||||
deadline := time.Now().Add(30 * time.Second)
|
||||
for {
|
||||
var missing []uint64
|
||||
for _, seq := range sent {
|
||||
if _, ok := got.Load(seq); !ok {
|
||||
missing = append(missing, seq)
|
||||
}
|
||||
}
|
||||
if len(missing) == 0 {
|
||||
return
|
||||
}
|
||||
info, err := js.ConsumerInfo(stream, "controller")
|
||||
if time.Now().After(deadline) || (err == nil && info.NumPending == 0 && info.NumAckPending == 0) {
|
||||
t.Fatalf("%d of %d followed events were never handed to the consumer (first at stream sequence %d), "+
|
||||
"and it has nothing pending: the server moved past them — issue 266, on %s", len(missing), len(sent),
|
||||
missing[0], admin.ConnectedServerVersion())
|
||||
}
|
||||
time.Sleep(200 * time.Millisecond)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,229 @@
|
||||
// prove runs each replay on the commit before its fix and on its fix, and says whether it fails before
|
||||
// and passes after (novox/hq to-be 45 Phase 5, "done when"). A replay that passes before its fix proves
|
||||
// nothing about it, and one that fails after it is not a replay of that fix.
|
||||
//
|
||||
// go run ./cmd/prove [--repos <dir holding the clones, ../.. by default>] [R262 R266 …]
|
||||
//
|
||||
// The clones are the core repositories, side by side as the lab expects them. A replay in its repository
|
||||
// is laid over a worktree of each commit and run with go test; the bus replay runs against a server of
|
||||
// the release the catalogue pinned at each commit, raised for it; the resolver replay with the catalogue's
|
||||
// resolver at each commit. MESH_TEST_POSTGRES is passed through to the controller's.
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"flag"
|
||||
"fmt"
|
||||
"net"
|
||||
"os"
|
||||
"os/exec"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/novox/mesh-lab/replays"
|
||||
)
|
||||
|
||||
func main() {
|
||||
repos := flag.String("repos", filepath.Join("..", ".."), "the directory the core repositories are cloned side by side in")
|
||||
flag.Parse()
|
||||
if abs, err := filepath.Abs(*repos); err == nil {
|
||||
*repos = abs
|
||||
}
|
||||
want := flag.Args()
|
||||
failed := 0
|
||||
for _, r := range replays.Register {
|
||||
if len(want) > 0 && !contains(want, r.ID) {
|
||||
continue
|
||||
}
|
||||
at := r.Fix + "^1"
|
||||
if r.Before != "" {
|
||||
at = r.Before
|
||||
}
|
||||
before, beforeSaid := run(r, *repos, at)
|
||||
after, afterSaid := run(r, *repos, r.Fix)
|
||||
verdict := "PROVED"
|
||||
if before != "fails" && before != "absent" || after != "passes" {
|
||||
verdict = "NOT PROVED"
|
||||
failed++
|
||||
}
|
||||
fmt.Printf("%s (issue %d): %s — before its fix %s, on its fix %s\n", r.ID, r.Issue, verdict, before, after)
|
||||
if verdict != "PROVED" {
|
||||
fmt.Printf(" before: %s\n after: %s\n", beforeSaid, afterSaid)
|
||||
} else {
|
||||
fmt.Printf(" before: %s\n", firstLine(beforeSaid))
|
||||
}
|
||||
}
|
||||
if failed > 0 {
|
||||
os.Exit(1)
|
||||
}
|
||||
}
|
||||
|
||||
// run is the replay at one commit: passes, fails, or absent (the check it replays did not exist).
|
||||
func run(r replays.Replay, repos, at string) (string, string) {
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 15*time.Minute)
|
||||
defer cancel()
|
||||
repo := filepath.Join(repos, r.Repository)
|
||||
tree, err := os.MkdirTemp("", "prove-"+r.ID+"-")
|
||||
if err != nil {
|
||||
return "error", err.Error()
|
||||
}
|
||||
defer os.RemoveAll(tree)
|
||||
if out, err := exec.CommandContext(ctx, "git", "-C", repo, "worktree", "add", "--quiet", "--detach", tree, at).CombinedOutput(); err != nil {
|
||||
return "error", string(out)
|
||||
}
|
||||
defer exec.Command("git", "-C", repo, "worktree", "remove", "--force", tree).Run()
|
||||
|
||||
switch r.Kind {
|
||||
case replays.InRepository, replays.Gate:
|
||||
home := filepath.Join(repos, r.Repository)
|
||||
if r.Home != "" {
|
||||
home = filepath.Join(repos, r.Home)
|
||||
}
|
||||
ref := r.HomeRef
|
||||
if ref == "" {
|
||||
ref = r.Fix
|
||||
}
|
||||
for _, f := range r.Files {
|
||||
body, err := exec.Command("git", "-C", home, "show", ref+":"+f).Output()
|
||||
if err != nil {
|
||||
return "error", fmt.Sprintf("%s has no %s: %v", ref, f, err)
|
||||
}
|
||||
if err := os.WriteFile(filepath.Join(tree, f), body, 0o644); err != nil {
|
||||
return "error", err.Error()
|
||||
}
|
||||
}
|
||||
mode := "-mod=vendor"
|
||||
if _, err := os.Stat(filepath.Join(tree, "vendor")); err != nil {
|
||||
mode = "-mod=mod"
|
||||
}
|
||||
cmd := exec.CommandContext(ctx, "go", "test", "-count=1", "-run", "^"+r.Test, r.Package)
|
||||
cmd.Dir = tree
|
||||
cmd.Env = append(os.Environ(), "GOFLAGS="+mode, "GOPRIVATE=git.novox.be")
|
||||
out, err := cmd.CombinedOutput()
|
||||
return verdictOf(string(out), err, r.Kind == replays.Gate)
|
||||
case replays.Bus:
|
||||
// **The very image the catalogue pinned at the commit**, by its digest: the release a digest is was
|
||||
// not said beside it before the fix, and the digest is what ran.
|
||||
image, err := pinnedBus(tree)
|
||||
if err != nil {
|
||||
return "error", err.Error()
|
||||
}
|
||||
url, stop, err := aBus(ctx, image)
|
||||
if err != nil {
|
||||
return "error", err.Error()
|
||||
}
|
||||
defer stop()
|
||||
return replayHere(ctx, r, "MESH_TEST_NATS="+url)
|
||||
case replays.Resolver:
|
||||
return replayHere(ctx, r, "MESH_REPLAY_CATALOGUE="+tree)
|
||||
}
|
||||
return "error", "no such kind"
|
||||
}
|
||||
|
||||
func replayHere(ctx context.Context, r replays.Replay, env string) (string, string) {
|
||||
cmd := exec.CommandContext(ctx, "go", "test", "-count=1", "-run", "^"+r.Test, r.Package)
|
||||
cmd.Env = append(os.Environ(), env)
|
||||
out, err := cmd.CombinedOutput()
|
||||
return verdictOf(string(out), err, false)
|
||||
}
|
||||
|
||||
// verdictOf reads go test's word: a replay that does not build where the check did not exist is absent.
|
||||
func verdictOf(out string, err error, gate bool) (string, string) {
|
||||
switch {
|
||||
case err == nil && strings.Contains(out, "no tests to run"):
|
||||
return "absent", "no test of that name at this commit"
|
||||
case err == nil:
|
||||
return "passes", lastLines(out, 3)
|
||||
case strings.Contains(out, "[build failed]") || strings.Contains(out, "[setup failed]"):
|
||||
if gate {
|
||||
return "absent", "the check did not exist: " + lastLines(out, 3)
|
||||
}
|
||||
return "error", "does not build: " + lastLines(out, 6)
|
||||
case strings.Contains(out, "--- FAIL"):
|
||||
return "fails", failure(out)
|
||||
}
|
||||
return "error", lastLines(out, 6)
|
||||
}
|
||||
|
||||
// pinnedBus is the bus image the catalogue's nats module is built on at a commit.
|
||||
func pinnedBus(tree string) (string, error) {
|
||||
raw, err := os.ReadFile(filepath.Join(tree, "modules", "nats", "module.json"))
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
var m struct {
|
||||
Build struct {
|
||||
On []struct{ Arg, Image string } `json:"on"`
|
||||
} `json:"build"`
|
||||
}
|
||||
if err := json.Unmarshal(raw, &m); err != nil {
|
||||
return "", err
|
||||
}
|
||||
for _, on := range m.Build.On {
|
||||
if on.Arg == "NATS_BASE" && on.Image != "" {
|
||||
return on.Image, nil
|
||||
}
|
||||
}
|
||||
return "", errors.New("the catalogue's nats module names no NATS_BASE image")
|
||||
}
|
||||
|
||||
// aBus is a server of an image, on loopback, for the length of one replay.
|
||||
func aBus(ctx context.Context, image string) (string, func(), error) {
|
||||
name := fmt.Sprintf("prove-bus-%d", time.Now().UnixNano())
|
||||
if out, err := exec.CommandContext(ctx, "docker", "run", "-d", "--rm", "--name", name, "-p", "127.0.0.1::4222",
|
||||
image, "-js").CombinedOutput(); err != nil {
|
||||
return "", nil, fmt.Errorf("%s: %v %s", image, err, out)
|
||||
}
|
||||
stop := func() { _ = exec.Command("docker", "rm", "-f", name).Run() }
|
||||
out, err := exec.CommandContext(ctx, "docker", "port", name, "4222/tcp").Output()
|
||||
if err != nil {
|
||||
stop()
|
||||
return "", nil, err
|
||||
}
|
||||
address := strings.TrimSpace(strings.Split(string(out), "\n")[0])
|
||||
for i := 0; i < 50; i++ {
|
||||
if c, err := net.Dial("tcp", address); err == nil {
|
||||
c.Close()
|
||||
break
|
||||
}
|
||||
time.Sleep(200 * time.Millisecond)
|
||||
}
|
||||
return "nats://" + address, stop, nil
|
||||
}
|
||||
|
||||
func failure(out string) string {
|
||||
for _, line := range strings.Split(out, "\n") {
|
||||
if strings.Contains(line, "_test.go:") {
|
||||
return strings.TrimSpace(line)
|
||||
}
|
||||
}
|
||||
return lastLines(out, 3)
|
||||
}
|
||||
|
||||
func lastLines(s string, n int) string {
|
||||
lines := strings.Split(strings.TrimSpace(s), "\n")
|
||||
if len(lines) > n {
|
||||
lines = lines[len(lines)-n:]
|
||||
}
|
||||
return strings.Join(lines, " | ")
|
||||
}
|
||||
|
||||
func firstLine(s string) string {
|
||||
line, _, _ := strings.Cut(s, "\n")
|
||||
if len(line) > 300 {
|
||||
line = line[:300] + "…"
|
||||
}
|
||||
return line
|
||||
}
|
||||
|
||||
func contains(list []string, s string) bool {
|
||||
for _, x := range list {
|
||||
if x == s {
|
||||
return true
|
||||
}
|
||||
}
|
||||
return false
|
||||
}
|
||||
@@ -0,0 +1,196 @@
|
||||
package replays
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"encoding/binary"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"io"
|
||||
"net"
|
||||
"net/http"
|
||||
"net/url"
|
||||
"os"
|
||||
"strings"
|
||||
"time"
|
||||
)
|
||||
|
||||
// Docker is the container runtime's own API, over its socket: what a replay needs to raise a resolver and
|
||||
// ask it from programs under two C libraries — pull, network, run, wait, logs, remove — and nothing else.
|
||||
// Spoken directly rather than through the command line, because the Go toolchain a check runs in carries
|
||||
// no docker client, and a replay that needed one would be a replay that runs only on a workstation.
|
||||
type Docker struct {
|
||||
http *http.Client
|
||||
// Label marks everything it makes, so a replay that dies is cleaned up by it.
|
||||
Label string
|
||||
}
|
||||
|
||||
// DockerFromEnv is the runtime at DOCKER_HOST's socket, or the usual one; ok false when there is none.
|
||||
func DockerFromEnv() (*Docker, bool) {
|
||||
sock := "/var/run/docker.sock"
|
||||
if h := os.Getenv("DOCKER_HOST"); strings.HasPrefix(h, "unix://") {
|
||||
sock = strings.TrimPrefix(h, "unix://")
|
||||
}
|
||||
if _, err := os.Stat(sock); err != nil {
|
||||
return nil, false
|
||||
}
|
||||
return &Docker{Label: "mesh.replay", http: &http.Client{Transport: &http.Transport{
|
||||
DialContext: func(ctx context.Context, _, _ string) (net.Conn, error) {
|
||||
return (&net.Dialer{}).DialContext(ctx, "unix", sock)
|
||||
}}}}, true
|
||||
}
|
||||
|
||||
func (d *Docker) call(ctx context.Context, method, path string, body any, out any) error {
|
||||
var r io.Reader
|
||||
if body != nil {
|
||||
raw, err := json.Marshal(body)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
r = bytes.NewReader(raw)
|
||||
}
|
||||
req, err := http.NewRequestWithContext(ctx, method, "http://docker"+path, r)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
req.Header.Set("Content-Type", "application/json")
|
||||
res, err := d.http.Do(req)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer res.Body.Close()
|
||||
raw, _ := io.ReadAll(res.Body)
|
||||
if res.StatusCode >= 300 {
|
||||
return fmt.Errorf("docker %s %s: %s %s", method, path, res.Status, strings.TrimSpace(string(raw)))
|
||||
}
|
||||
if out != nil && len(raw) > 0 {
|
||||
return json.Unmarshal(raw, out)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// Pull fetches an image unless the runtime has it.
|
||||
func (d *Docker) Pull(ctx context.Context, image string) error {
|
||||
if err := d.call(ctx, http.MethodGet, "/images/"+url.PathEscape(image)+"/json", nil, nil); err == nil {
|
||||
return nil
|
||||
}
|
||||
name, tag, _ := strings.Cut(image, ":")
|
||||
if tag == "" {
|
||||
tag = "latest"
|
||||
}
|
||||
return d.call(ctx, http.MethodPost, "/images/create?fromImage="+url.QueryEscape(name)+"&tag="+url.QueryEscape(tag), nil, nil)
|
||||
}
|
||||
|
||||
// Network makes a network of its own and answers its id.
|
||||
func (d *Docker) Network(ctx context.Context, name string) (string, error) {
|
||||
var made struct{ ID string }
|
||||
err := d.call(ctx, http.MethodPost, "/networks/create", map[string]any{"Name": name,
|
||||
"Labels": map[string]string{d.Label: name}}, &made)
|
||||
return made.ID, err
|
||||
}
|
||||
|
||||
// RemoveNetwork removes a network.
|
||||
func (d *Docker) RemoveNetwork(ctx context.Context, id string) {
|
||||
_ = d.call(ctx, http.MethodDelete, "/networks/"+id, nil, nil)
|
||||
}
|
||||
|
||||
// Run is one container: an image, a command, the files it is given, its network and resolver.
|
||||
type Run struct {
|
||||
Image string
|
||||
Cmd []string
|
||||
Binds []string
|
||||
Network string
|
||||
DNS []string
|
||||
}
|
||||
|
||||
// Start creates and starts a container and answers its id and its address on its network.
|
||||
func (d *Docker) Start(ctx context.Context, r Run) (string, string, error) {
|
||||
var made struct{ ID string }
|
||||
host := map[string]any{"Binds": r.Binds, "NetworkMode": r.Network}
|
||||
if len(r.DNS) > 0 {
|
||||
host["Dns"] = r.DNS
|
||||
}
|
||||
if err := d.call(ctx, http.MethodPost, "/containers/create", map[string]any{"Image": r.Image, "Cmd": r.Cmd,
|
||||
"Labels": map[string]string{d.Label: "1"}, "HostConfig": host}, &made); err != nil {
|
||||
return "", "", err
|
||||
}
|
||||
if err := d.call(ctx, http.MethodPost, "/containers/"+made.ID+"/start", nil, nil); err != nil {
|
||||
d.Remove(ctx, made.ID)
|
||||
return "", "", err
|
||||
}
|
||||
var seen struct {
|
||||
NetworkSettings struct {
|
||||
Networks map[string]struct{ IPAddress string }
|
||||
}
|
||||
}
|
||||
if err := d.call(ctx, http.MethodGet, "/containers/"+made.ID+"/json", nil, &seen); err != nil {
|
||||
return made.ID, "", err
|
||||
}
|
||||
for _, n := range seen.NetworkSettings.Networks {
|
||||
return made.ID, n.IPAddress, nil
|
||||
}
|
||||
return made.ID, "", nil
|
||||
}
|
||||
|
||||
// Wait waits for a container to end, and answers its exit code and what it said.
|
||||
func (d *Docker) Wait(ctx context.Context, id string) (int, string, error) {
|
||||
var ended struct{ StatusCode int }
|
||||
if err := d.call(ctx, http.MethodPost, "/containers/"+id+"/wait", nil, &ended); err != nil {
|
||||
return -1, "", err
|
||||
}
|
||||
return ended.StatusCode, d.Logs(ctx, id), nil
|
||||
}
|
||||
|
||||
// Logs is what a container said, both streams.
|
||||
func (d *Docker) Logs(ctx context.Context, id string) string {
|
||||
req, err := http.NewRequestWithContext(ctx, http.MethodGet, "http://docker/containers/"+id+"/logs?stdout=1&stderr=1", nil)
|
||||
if err != nil {
|
||||
return ""
|
||||
}
|
||||
res, err := d.http.Do(req)
|
||||
if err != nil {
|
||||
return ""
|
||||
}
|
||||
defer res.Body.Close()
|
||||
var b strings.Builder
|
||||
head := make([]byte, 8)
|
||||
for {
|
||||
if _, err := io.ReadFull(res.Body, head); err != nil {
|
||||
return b.String()
|
||||
}
|
||||
n := binary.BigEndian.Uint32(head[4:])
|
||||
if _, err := io.CopyN(&b, res.Body, int64(n)); err != nil {
|
||||
return b.String()
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Remove removes a container, running or not.
|
||||
func (d *Docker) Remove(ctx context.Context, id string) {
|
||||
_ = d.call(ctx, http.MethodDelete, "/containers/"+id+"?force=1", nil, nil)
|
||||
}
|
||||
|
||||
// Exec runs a command in a running container and answers its exit code.
|
||||
func (d *Docker) Exec(ctx context.Context, id string, cmd ...string) (int, error) {
|
||||
var made struct{ ID string }
|
||||
if err := d.call(ctx, http.MethodPost, "/containers/"+id+"/exec", map[string]any{"Cmd": cmd}, &made); err != nil {
|
||||
return -1, err
|
||||
}
|
||||
if err := d.call(ctx, http.MethodPost, "/exec/"+made.ID+"/start", map[string]any{"Detach": true}, nil); err != nil {
|
||||
return -1, err
|
||||
}
|
||||
for i := 0; i < 600; i++ {
|
||||
var seen struct {
|
||||
Running bool
|
||||
ExitCode int
|
||||
}
|
||||
if err := d.call(ctx, http.MethodGet, "/exec/"+made.ID+"/json", nil, &seen); err != nil {
|
||||
return -1, err
|
||||
}
|
||||
if !seen.Running {
|
||||
return seen.ExitCode, nil
|
||||
}
|
||||
time.Sleep(100 * time.Millisecond)
|
||||
}
|
||||
return -1, fmt.Errorf("%v did not end", cmd)
|
||||
}
|
||||
@@ -0,0 +1,13 @@
|
||||
module github.com/novox/mesh-lab/replays
|
||||
|
||||
go 1.26.0
|
||||
|
||||
require github.com/nats-io/nats.go v1.54.0
|
||||
|
||||
require (
|
||||
github.com/klauspost/compress v1.20.0 // 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
|
||||
)
|
||||
@@ -0,0 +1,12 @@
|
||||
github.com/klauspost/compress v1.20.0 h1:a3C1ke2ohxFymNlb2HWAHjDeKCI90scRskErZkR0ezA=
|
||||
github.com/klauspost/compress v1.20.0/go.mod h1:LUdAzn7YLVvxLpc7y3V1m40wESHTgc1422pwwBSKYuI=
|
||||
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.48.0 h1:bbX/i/6MgT9BVLM9RT1thmxL04yeTAhbEz4SyadbXoo=
|
||||
golang.org/x/sys v0.48.0/go.mod h1:hNLxWAXmnKAxqDtdwIYC4bM9oQPEecfsnNMuSxOs3og=
|
||||
@@ -0,0 +1,91 @@
|
||||
// Package replays is the lab's replay of every core incident (novox/hq to-be 45 §9, M9): a scripted
|
||||
// replay of what happened, asserting the rule's outcome, run before every merge of a core repository —
|
||||
// and proved, once, to fail on the commit before its fix and to pass on the fix.
|
||||
//
|
||||
// **Where a replay lives.** A replay of one component's logic is a test in that component's repository,
|
||||
// written only with what the component had before its fix, so it can be laid over the older commit
|
||||
// (the controller's replays_test.go). A replay of what the mesh runs rather than what it wrote — the bus
|
||||
// server's release, the resolver the catalogue configures, under both C libraries — lives here, where a
|
||||
// container runtime and a bus of a given release are at hand. Every one is in Register, with the issue,
|
||||
// the fix, and how the prover runs it at a commit.
|
||||
//
|
||||
// **A new core issue resolves with its replay added here, or with a stated reason none is possible**
|
||||
// (hq 00-META/checks/cycle.py holds the issue to it: `replay:` names an entry here, or `replay-none:`
|
||||
// says why).
|
||||
package replays
|
||||
|
||||
// Kind is how a replay is run at a commit.
|
||||
type Kind string
|
||||
|
||||
const (
|
||||
// InRepository is a test in the fixed repository, laid over the commit and run there with go test.
|
||||
InRepository Kind = "test"
|
||||
// Bus is the bus replay here, run against a server of the release the catalogue pinned at the commit.
|
||||
Bus Kind = "bus"
|
||||
// Resolver is the resolver replay here, run with the catalogue's resolver configuration at the commit,
|
||||
// asked by a program under musl and one under glibc.
|
||||
Resolver Kind = "resolver"
|
||||
// Gate is a test of the merge gate itself: before its fix there was no check to fail, so the commit
|
||||
// before is held to fail by not having it.
|
||||
Gate Kind = "gate"
|
||||
)
|
||||
|
||||
// Replay is one incident, replayed.
|
||||
type Replay struct {
|
||||
ID string
|
||||
Issue int
|
||||
// What is the incident in a line, and Asserts the rule's outcome the replay holds.
|
||||
What, Asserts string
|
||||
Kind Kind
|
||||
// Repository is where the fix is, and Fix the fix's commit there: the replay must fail on Fix's first
|
||||
// parent and pass on Fix.
|
||||
Repository string
|
||||
Fix string
|
||||
// Before is the commit the replay must fail on, when that is not Fix's first parent: a fix still on its
|
||||
// branch is proved against the main it branched from. Set to "" once the fix is merged.
|
||||
Before string
|
||||
// Package and Test are where the replay is run: a package of Repository for InRepository and Gate,
|
||||
// of this module otherwise; Files are the files laid over the commit, from where the replay lives.
|
||||
Package string
|
||||
Test string
|
||||
Files []string
|
||||
// Home is where the replay's files live, when that is not the fix: the repository and ref.
|
||||
Home, HomeRef string
|
||||
}
|
||||
|
||||
// Register is every replay, by incident.
|
||||
var Register = []Replay{
|
||||
{ID: "R236", Issue: 236, Kind: Gate, Repository: "mesh-controller", Fix: "feat/merge-gate",
|
||||
Before: "origin/main",
|
||||
What: "a login manager's service passed the catalogue check and was refused whole by the node-engine",
|
||||
Asserts: "a manifest the node-engine refuses fails its pull request, naming the machine and the module",
|
||||
Package: "./cmd/mesh-controller", Test: "TestIssue236", Files: []string{"cmd/mesh-controller/merge_gate_test.go"}},
|
||||
{ID: "R262", Issue: 262, Kind: Resolver, Repository: "mesh-catalog", Fix: "20603b63e67323f2c725b242c0ea87281d4dee93",
|
||||
What: "an Alpine container could not find a machine by its mesh name: NXDOMAIN for IPv6 is final to musl",
|
||||
Asserts: "every machine's name resolves through the catalogue's resolver under musl and under glibc",
|
||||
Package: ".", Test: "TestReplay262"},
|
||||
{ID: "R263", Issue: 263, Kind: InRepository, Repository: "mesh-controller", Fix: "6d620f77c3f3f967be953ff760823c6a51b7b8d3",
|
||||
What: "a consumer's 26-character identity refused the resolver holder's whole declaration",
|
||||
Asserts: "a provider's machine composes whatever its consumers are called",
|
||||
Package: "./cmd/mesh-controller", Test: "TestReplay263", Files: []string{"cmd/mesh-controller/replays_test.go"},
|
||||
Home: "mesh-controller", HomeRef: "feat/replays"},
|
||||
{ID: "R266", Issue: 266, Kind: Bus, Repository: "mesh-catalog", Fix: "0275c2e",
|
||||
What: "a merge on the bus was never handed to the controller: 2.10 skipped messages on a consumer with several filters",
|
||||
Asserts: "a consumer filtered like the controller's is handed every followed message under the mesh's traffic",
|
||||
Package: ".", Test: "TestReplay266"},
|
||||
{ID: "R273", Issue: 273, Kind: InRepository, Repository: "mesh-controller", Fix: "8bfaf1523ed2dddde80daa0994f1044e59db44b1",
|
||||
What: "a rule for the resolver re-bound a machine's database consumers to empty databases elsewhere",
|
||||
Asserts: "a consumer beside its store stays bound to it while another machine holds the store's seat",
|
||||
Package: "./cmd/mesh-controller", Test: "TestReplay273", Files: []string{"cmd/mesh-controller/replays_test.go"},
|
||||
Home: "mesh-controller", HomeRef: "feat/replays"},
|
||||
}
|
||||
|
||||
// Find is the replay of that id, or of that issue.
|
||||
func Find(id string) (Replay, bool) {
|
||||
for _, r := range Register {
|
||||
if r.ID == id || r.ID == "R"+id {
|
||||
return r, true
|
||||
}
|
||||
}
|
||||
return Replay{}, false
|
||||
}
|
||||
@@ -0,0 +1,37 @@
|
||||
package replays
|
||||
|
||||
import (
|
||||
"regexp"
|
||||
"testing"
|
||||
)
|
||||
|
||||
// Every replay says what happened, what it asserts, where its fix is, and how it is run.
|
||||
func TestEveryReplayIsWhole(t *testing.T) {
|
||||
seen := map[string]bool{}
|
||||
for _, r := range Register {
|
||||
if seen[r.ID] {
|
||||
t.Errorf("%s is registered twice", r.ID)
|
||||
}
|
||||
seen[r.ID] = true
|
||||
if r.Issue == 0 || r.What == "" || r.Asserts == "" || r.Repository == "" || r.Fix == "" || r.Test == "" {
|
||||
t.Errorf("%s does not say its incident, its outcome, its fix and its test: %+v", r.ID, r)
|
||||
}
|
||||
if !regexp.MustCompile(`^R\d+$`).MatchString(r.ID) {
|
||||
t.Errorf("%q is not a replay's id", r.ID)
|
||||
}
|
||||
switch r.Kind {
|
||||
case InRepository, Gate:
|
||||
if len(r.Files) == 0 || r.Package == "" {
|
||||
t.Errorf("%s runs in its repository and lays no files over the commit", r.ID)
|
||||
}
|
||||
case Bus, Resolver:
|
||||
default:
|
||||
t.Errorf("%s is of no kind the prover runs: %q", r.ID, r.Kind)
|
||||
}
|
||||
}
|
||||
for _, issue := range []string{"236", "262", "263", "266"} {
|
||||
if _, ok := Find(issue); !ok {
|
||||
t.Errorf("to-be 45 Phase 5 is done when the replay of %s fails before its fix and passes after; it has none", issue)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,133 @@
|
||||
package replays
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"encoding/base64"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"testing"
|
||||
"text/template"
|
||||
"time"
|
||||
)
|
||||
|
||||
// **R262 — every machine's name resolves under musl and under glibc** (novox/hq issue 262; to-be 45 §9,
|
||||
// "the resolver module's tests run under musl and glibc").
|
||||
//
|
||||
// After the mesh moved to one resolver, an Alpine container crash-looped on `getaddrinfo ENOTFOUND` for
|
||||
// a machine's mesh name that `getent hosts` found: the resolver answered the name only through a
|
||||
// wildcard, so asked for its IPv6 address it said there is no such name, and musl — which asks for both
|
||||
// and takes NXDOMAIN as final — gave up; glibc did not. The replay renders the machine list the
|
||||
// catalogue's resolver module reads (its `node-zones` fact, from the manifest at MESH_REPLAY_CATALOGUE),
|
||||
// runs the resolver with it, and asks for each machine by `getaddrinfo` from a program under each C
|
||||
// library, on the runtime's default bridge as the incident's container was. Needs a container runtime
|
||||
// (its socket); its images come from the public registry.
|
||||
var libcs = map[string]string{"musl": "python:3.13-alpine", "glibc": "python:3.13-slim"}
|
||||
|
||||
// resolverImage is where the resolver runs: the distribution's own dnsmasq, installed as the module
|
||||
// installs it on a machine.
|
||||
const resolverImage = "alpine:3.20"
|
||||
|
||||
func TestReplay262EveryMachineResolvesUnderMuslAndGlibc(t *testing.T) {
|
||||
catalogue := os.Getenv("MESH_REPLAY_CATALOGUE")
|
||||
if catalogue == "" {
|
||||
catalogue = filepath.Join("..", "..", "mesh-catalog")
|
||||
}
|
||||
raw, err := os.ReadFile(filepath.Join(catalogue, "modules", "dnsmasq", "module.json"))
|
||||
if err != nil {
|
||||
t.Skipf("no catalogue to read the resolver's configuration from (%v); MESH_REPLAY_CATALOGUE names one", err)
|
||||
}
|
||||
docker, ok := DockerFromEnv()
|
||||
if !ok {
|
||||
t.Skip("no container runtime: the resolver and the two C libraries run in containers")
|
||||
}
|
||||
var manifest struct {
|
||||
Facts map[string]struct {
|
||||
Template string `json:"template"`
|
||||
} `json:"facts"`
|
||||
}
|
||||
if err := json.Unmarshal(raw, &manifest); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
zones, ok := manifest.Facts["node-zones"]
|
||||
if !ok {
|
||||
t.Fatal("the resolver module renders no machine list (`facts.node-zones`)")
|
||||
}
|
||||
type entry struct{ Name, FQDN, Address, Account string }
|
||||
machines := []entry{{Name: "anchor", FQDN: "anchor.internal", Address: "10.88.0.1"},
|
||||
{Name: "homeserver", FQDN: "homeserver.internal", Address: "10.88.0.2"}}
|
||||
var rendered bytes.Buffer
|
||||
tpl, err := template.New("node-zones").Parse(zones.Template)
|
||||
if err == nil {
|
||||
err = tpl.Execute(&rendered, map[string]any{"Node": "anchor", "Suffix": "internal", "Names": machines,
|
||||
"Machines": machines, "Zones": nil, "Holders": map[string][]entry{}})
|
||||
}
|
||||
if err != nil {
|
||||
t.Fatalf("the machine list does not render: %v", err)
|
||||
}
|
||||
|
||||
ctx, cancel := context.WithTimeout(t.Context(), 5*time.Minute)
|
||||
defer cancel()
|
||||
for _, image := range append([]string{resolverImage}, libcs["musl"], libcs["glibc"]) {
|
||||
if err := docker.Pull(ctx, image); err != nil {
|
||||
t.Fatalf("cannot fetch %s: %v", image, err)
|
||||
}
|
||||
}
|
||||
// The resolver: nothing but the machine list, answering on its address on the bridge.
|
||||
conf := base64.StdEncoding.EncodeToString(rendered.Bytes())
|
||||
script := fmt.Sprintf(`apk add -q --no-cache dnsmasq >/dev/null && echo %s | base64 -d > /etc/nodes.conf && `+
|
||||
`exec dnsmasq -k --no-resolv --no-hosts --log-queries --log-facility=- --conf-file=/etc/nodes.conf`, conf)
|
||||
resolver, address, err := docker.Start(ctx, Run{Image: resolverImage, Cmd: []string{"sh", "-c", script}, Network: "bridge"})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer docker.Remove(context.Background(), resolver)
|
||||
if address == "" {
|
||||
t.Fatal("the resolver has no address on the bridge")
|
||||
}
|
||||
|
||||
// Each machine, asked by getaddrinfo as a program does — both families, as musl asks — a few times
|
||||
// while the resolver comes up.
|
||||
ask := fmt.Sprintf(`
|
||||
import socket, sys, time
|
||||
names = %q.split()
|
||||
for attempt in range(30):
|
||||
failed = []
|
||||
for n in names:
|
||||
try:
|
||||
socket.getaddrinfo(n, 80)
|
||||
except socket.gaierror as e:
|
||||
failed.append("%%s: %%s" %% (n, e))
|
||||
if not failed:
|
||||
print("resolved", " ".join(names)); sys.exit(0)
|
||||
time.sleep(1)
|
||||
print("\n".join(failed)); sys.exit(1)
|
||||
`, machines[0].FQDN+" "+machines[1].FQDN)
|
||||
for libc, image := range libcs {
|
||||
id, _, err := docker.Start(ctx, Run{Image: image, Cmd: []string{"python3", "-c", ask}, Network: "bridge",
|
||||
DNS: []string{address}})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
code, said, err := docker.Wait(ctx, id)
|
||||
docker.Remove(context.Background(), id)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if code != 0 {
|
||||
t.Errorf("under %s (%s) a machine's mesh name does not resolve through the catalogue's resolver — "+
|
||||
"issue 262:\n%s\nthe resolver said:\n%s", libc, image, strings.TrimSpace(said), tail(docker.Logs(ctx, resolver), 12))
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func tail(s string, n int) string {
|
||||
lines := strings.Split(strings.TrimSpace(s), "\n")
|
||||
if len(lines) > n {
|
||||
lines = lines[len(lines)-n:]
|
||||
}
|
||||
return strings.Join(lines, "\n")
|
||||
}
|
||||
Reference in New Issue
Block a user