Compare commits
5
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
0a2e58c070 | ||
|
|
45b4ae92bc | ||
|
|
ed45cc6415 | ||
|
|
ef551fdfb6 | ||
|
|
f510b46319 |
@@ -5,8 +5,6 @@ import (
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"path"
|
||||
"regexp"
|
||||
"runtime/debug"
|
||||
"slices"
|
||||
"sort"
|
||||
"strings"
|
||||
@@ -248,11 +246,9 @@ func checkRequestFor(ctx context.Context, open *stores, p link.PullUpdated, scop
|
||||
if dir == "mesh-controller" {
|
||||
beside["mesh-controller-main"] = link.CheckedOut{Repository: url, Ref: refs["mesh-controller-main"]}
|
||||
if e.Source.Seat != "" {
|
||||
for _, sibling := range []string{"mesh-lab", "mesh-sdk"} {
|
||||
if url, err := clone(inventory.Source{Seat: e.Source.Seat, Repository: siblingOf(e.Source.Repository,
|
||||
sibling)}); err == nil {
|
||||
beside[sibling] = link.CheckedOut{Repository: url, Ref: refs[sibling]}
|
||||
}
|
||||
if lab, err := clone(inventory.Source{Seat: e.Source.Seat, Repository: siblingOf(e.Source.Repository,
|
||||
"mesh-lab")}); err == nil {
|
||||
beside["mesh-lab"] = link.CheckedOut{Repository: lab, Ref: refs["mesh-lab"]}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -275,54 +271,17 @@ func checkRequestFor(ctx context.Context, open *stores, p link.PullUpdated, scop
|
||||
// catalogue, whose checkout beside is what tests read its files from, at its main, what the next merge
|
||||
// builds from; and beside the controller its main, for a judge the running controller predates, and the
|
||||
// lab's main, whose replays every check runs. **One rule, read by the check the controller asks for and by
|
||||
// the facts snapshot** (Facts.Beside), so a check run by hand clones what the build seat clones. Beside the
|
||||
// controller also the SDK, checked out at the commit the running controller's go.mod pins (sdkPinned), not
|
||||
// one a desktop holds (novox/hq issue 449). The checkout only places the clone: the conformance test reads the
|
||||
// fixtures at the pin of the tree under check, from the clone's history, so a pull request moving the SDK is
|
||||
// judged against the SDK it moves to (internal/link/conformance_test.go).
|
||||
// the facts snapshot** (Facts.Beside), so a check run by hand clones what the build seat clones.
|
||||
func besideRefs(dir, running string) map[string]string {
|
||||
switch dir {
|
||||
case "mesh-catalog":
|
||||
return map[string]string{dir: "main"}
|
||||
case "mesh-controller":
|
||||
return map[string]string{dir: running, "mesh-controller-main": "main", "mesh-lab": "main",
|
||||
"mesh-sdk": sdkPinned()}
|
||||
return map[string]string{dir: running, "mesh-controller-main": "main", "mesh-lab": "main"}
|
||||
}
|
||||
return map[string]string{dir: running}
|
||||
}
|
||||
|
||||
// sdkModule is the Go module of the SDK the controller is built against.
|
||||
const sdkModule = "git.novox.be/novox/mesh-sdk/go"
|
||||
|
||||
// pseudoCommit is the commit a Go pseudo-version names: v0.1.11-0.20261009143344-f047d0a4a970 → f047d0a4a970.
|
||||
var pseudoCommit = regexp.MustCompile(`-([0-9a-f]{12})$`)
|
||||
|
||||
// sdkPinned is the ref of the SDK repository this controller was built from, as its go.mod pins it and its
|
||||
// build records it: a pseudo-version's commit, or a release's tag (the SDK tags its Go module under go/).
|
||||
// "main" only for a binary that records no SDK version — a local replace — which a running controller is
|
||||
// not (novox/hq issue 449).
|
||||
func sdkPinned() string {
|
||||
info, ok := debug.ReadBuildInfo()
|
||||
if !ok {
|
||||
return "main"
|
||||
}
|
||||
for _, dep := range info.Deps {
|
||||
if dep.Path != sdkModule {
|
||||
continue
|
||||
}
|
||||
if dep.Replace != nil {
|
||||
dep = dep.Replace
|
||||
}
|
||||
if m := pseudoCommit.FindStringSubmatch(dep.Version); m != nil {
|
||||
return m[1]
|
||||
}
|
||||
if strings.HasPrefix(dep.Version, "v") {
|
||||
return "go/" + dep.Version
|
||||
}
|
||||
}
|
||||
return "main"
|
||||
}
|
||||
|
||||
// siblingOf is another repository of the same owner: novox/mesh-controller → novox/mesh-lab.
|
||||
func siblingOf(repository, name string) string {
|
||||
if cut := strings.LastIndex(repository, "/"); cut >= 0 {
|
||||
|
||||
@@ -144,11 +144,15 @@ func actOnDeadLetter(ctx context.Context, on *busHandles, act string, id uint64,
|
||||
}
|
||||
switch act {
|
||||
case "deliver":
|
||||
_, to, err := link.DeliverAgain(on.js, id)
|
||||
delivered, to, err := link.DeliverAgain(on.js, id)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
answer["delivered_on"] = to
|
||||
if delivered.Original != "" {
|
||||
// What became of the ask in its seat's queue (novox/hq issue 334): removed, or left, and why.
|
||||
answer["original"] = delivered.Original
|
||||
}
|
||||
answer["done"] = fmt.Sprintf("dead letter %d was delivered again to %s, and nobody else; it is no longer kept",
|
||||
id, consumerWho(d.Stream, d.Consumer))
|
||||
case "drop":
|
||||
|
||||
@@ -0,0 +1,66 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/novox/mesh-controller/internal/conditions"
|
||||
"github.com/novox/mesh-controller/internal/inventory"
|
||||
"github.com/novox/mesh-controller/internal/link"
|
||||
)
|
||||
|
||||
// A provider whose wait fails the controller's check lists who waits on it (novox/hq issue 450).
|
||||
//
|
||||
// A wait that does not check out is judged unhealthy (ADR 0283 decision 3), so the provider raises its unhealthy
|
||||
// condition and its consumers are held under it (ADR 0240 rule 5). What is stored is what the machine said, the
|
||||
// wait unchecked: read as said, the provider seems to wait for the operator, and the held consumers were never
|
||||
// listed on the condition that is open.
|
||||
|
||||
// refusedWait is a wait for a secret the database's manifest does not declare: it fails the check.
|
||||
var refusedWait = inventory.Wait{Part: "postgres-database", Secret: "certificate", What: "the database's certificate"}
|
||||
|
||||
// judgedRefused is the provider's resource as the controller judges it: unhealthy, saying why.
|
||||
func judgedRefused() inventory.ResourceHealth {
|
||||
r := waitingDatabaseFor(refusedWait)
|
||||
r.State, r.Waits = link.StateUnhealthy, nil
|
||||
r.Reason = "says it waits for the secret certificate, which db does not declare"
|
||||
return r
|
||||
}
|
||||
|
||||
// The consumer's statement arrives after the provider's, so it is the consumer's judging (sayWaiters) that must list
|
||||
// it at the provider's unhealthy condition, made urgent by who waits on it.
|
||||
func TestAProviderWhoseWaitFailedItsCheckListsWhoWaitsOnIt(t *testing.T) {
|
||||
k, _ := withConditionsInMemory(t)
|
||||
stored := waitingDatabaseFor(refusedWait)
|
||||
open := judgeBoth(t, k, []inventory.ResourceHealth{judgedRefused()},
|
||||
map[string]inventory.NodeHealth{"anchor": {Node: "anchor", Resources: []inventory.ResourceHealth{stored}}},
|
||||
shopFailingBeside(stored))
|
||||
provider := conditionOf(open, moduleUnhealthyKey("db", "anchor"))
|
||||
if len(open) != 1 || provider == nil {
|
||||
t.Fatalf("a provider whose wait failed its check and a consumer of it raised %v; want the provider's "+
|
||||
"unhealthy alone", openKeysOf(open))
|
||||
}
|
||||
if said := provider.Evidence[0].Said; !strings.Contains(said, "shop on laptop") {
|
||||
t.Fatalf("the provider's unhealthy condition does not list shop on laptop as waiting on it: %s", said)
|
||||
}
|
||||
if provider.Severity != conditions.Urgent {
|
||||
t.Fatalf("the provider's unhealthy condition is %s with a consumer waiting on it; want urgent", provider.Severity)
|
||||
}
|
||||
}
|
||||
|
||||
// A provider that waits for a secret already given is unhealthy to its consumers too, before its own condition opens:
|
||||
// they are held under it as under any unhealthy provider, never as waiting for the operator, and its condition, when
|
||||
// open, is said as unhealthy with who waits on it.
|
||||
func TestAProviderWaitingForASecretAlreadyGivenIsUnhealthyToItsConsumers(t *testing.T) {
|
||||
given := holdingOf(nil, shopFailingBeside(waitingDatabase()), shopOnTheDatabase)
|
||||
given.waits.given["db@anchor"] = map[string]time.Time{"licence": time.Now().Add(-time.Hour)}
|
||||
by, held := given.heldWith("laptop", "shop", []inventory.ResourceHealth{failingConsumer("shop")})
|
||||
if !held || by.provider != theDatabase || len(by.waits) > 0 {
|
||||
t.Fatalf("shop under a database waiting for a licence already given: held %v under %v with waits %v; want "+
|
||||
"held under db on anchor as unhealthy, no waits", held, by.provider, by.waits)
|
||||
}
|
||||
if _, uncovered := given.waitingUncovered("laptop", "shop", []inventory.ResourceHealth{failingConsumer("shop")}); uncovered {
|
||||
t.Fatalf("a provider whose wait failed its check is said as waiting for another part")
|
||||
}
|
||||
}
|
||||
@@ -390,10 +390,9 @@ func TestAWithheldPathIsStoodInForByAPath(t *testing.T) {
|
||||
// controller's ask and by the facts — and finds the repository it checks from its origin.
|
||||
func TestACheckByHandClonesWhatTheSeatClones(t *testing.T) {
|
||||
for dir, refs := range map[string]map[string]string{
|
||||
"mesh-catalog": {"mesh-catalog": "main"},
|
||||
"mesh-host": {"mesh-host": "c0ffee"},
|
||||
"mesh-controller": {"mesh-controller": "c0ffee", "mesh-controller-main": "main", "mesh-lab": "main",
|
||||
"mesh-sdk": sdkPinnedByGoMod(t)},
|
||||
"mesh-catalog": {"mesh-catalog": "main"},
|
||||
"mesh-host": {"mesh-host": "c0ffee"},
|
||||
"mesh-controller": {"mesh-controller": "c0ffee", "mesh-controller-main": "main", "mesh-lab": "main"},
|
||||
} {
|
||||
got := besideRefs(dir, "c0ffee")
|
||||
for d, ref := range refs {
|
||||
@@ -413,24 +412,3 @@ func TestACheckByHandClonesWhatTheSeatClones(t *testing.T) {
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// sdkPinnedByGoMod is the SDK commit this tree's go.mod pins, read from the file rather than the build, so
|
||||
// sdkPinned is held to what the repository says (novox/hq issue 449).
|
||||
func sdkPinnedByGoMod(t *testing.T) string {
|
||||
t.Helper()
|
||||
raw, err := os.ReadFile(filepath.Join("..", "..", "go.mod"))
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
for _, line := range strings.Split(string(raw), "\n") {
|
||||
fields := strings.Fields(line)
|
||||
if len(fields) >= 2 && fields[0] == sdkModule {
|
||||
if m := pseudoCommit.FindStringSubmatch(fields[1]); m != nil {
|
||||
return m[1]
|
||||
}
|
||||
return "go/" + fields[1]
|
||||
}
|
||||
}
|
||||
t.Fatalf("go.mod requires no %s", sdkModule)
|
||||
return ""
|
||||
}
|
||||
|
||||
@@ -326,8 +326,10 @@ func judgeModuleHealth(ctx context.Context, inv *inventory.Inventory, k *conditi
|
||||
// sayWaiters observes a provider's open condition again, with who waits on it, from its machine's newest
|
||||
// statement. Nothing when its condition is not open: it is raised by its own statements, on its own looks.
|
||||
func sayWaiters(ctx context.Context, k *conditions.Keeper, hold *holding, p catalogue.Chosen, now time.Time) error {
|
||||
// Its statement as it was judged, its waits checked (novox/hq issue 450): a wait that failed the check is
|
||||
// unhealthy, so the condition built here is the one its own statement raised, and who waits on it is listed.
|
||||
var rs []inventory.ResourceHealth
|
||||
for _, r := range hold.healths[p.Node].Resources {
|
||||
for _, r := range hold.checked(p.Node) {
|
||||
if r.Module == p.Module && (r.State == link.StateUnhealthy || r.State == link.StateWaiting) {
|
||||
rs = append(rs, r)
|
||||
}
|
||||
|
||||
@@ -37,6 +37,32 @@ type holding struct {
|
||||
shelf map[string]catalogue.Manifest
|
||||
// providers memoises providerFor by machine, consumer and provision.
|
||||
providers map[string]providerLookup
|
||||
// waits is what the waits of a statement are checked against, read as they are asked for (novox/hq issue 450).
|
||||
waits operatorWaitFacts
|
||||
}
|
||||
|
||||
// checked is a machine's newest statement as the controller judges it: each wait checked (ADR 0283 decision 3), so
|
||||
// a wait that does not check out reads as unhealthy, as it did when the statement was judged (novox/hq issue 450).
|
||||
// What is stored is what the machine said, the waits unchecked; read as stored, a provider whose wait failed its
|
||||
// check seems to wait for the operator while its unhealthy condition is open.
|
||||
func (h *holding) checked(machine string) []inventory.ResourceHealth {
|
||||
rs := h.healths[machine].Resources
|
||||
mods := waitingModules(rs)
|
||||
if len(mods) == 0 {
|
||||
return rs
|
||||
}
|
||||
if h.waits.manifests == nil {
|
||||
h.waits.manifests = map[string]catalogue.Manifest{}
|
||||
}
|
||||
for _, module := range mods {
|
||||
if _, has := h.waits.manifests[module]; !has {
|
||||
if m, ok := h.manifestOf(module); ok {
|
||||
h.waits.manifests[module] = m
|
||||
}
|
||||
}
|
||||
}
|
||||
readWaitFacts(h.ctx, h.inv, machine, mods, nil, &h.waits)
|
||||
return checkWaiting(machine, rs, h.waits)
|
||||
}
|
||||
|
||||
type providerLookup struct {
|
||||
@@ -125,7 +151,7 @@ func (h *holding) providerState(p catalogue.Chosen) (unhealthy, waiting bool, wa
|
||||
return true, false, nil
|
||||
}
|
||||
}
|
||||
for _, r := range h.healths[p.Node].Resources {
|
||||
for _, r := range h.checked(p.Node) {
|
||||
if r.Module != p.Module {
|
||||
continue
|
||||
}
|
||||
|
||||
@@ -41,12 +41,21 @@ func failingConsumer(module string) inventory.ResourceHealth {
|
||||
State: link.StateUnhealthy, Reason: "http /health on web: answered 500", Check: "http", Needs: "postgres-database"}
|
||||
}
|
||||
|
||||
// dbSecrets is what the database's waits name: own secrets issued outside the mesh, so a wait for one checks out
|
||||
// while nobody gave it (ADR 0283 decision 3, novox/hq issue 450).
|
||||
var dbSecrets = catalogue.OwnSecrets{
|
||||
"licence": {Path: "/s/licence", IssuedBy: catalogue.IssuedOutside},
|
||||
"backup-key": {Path: "/s/backup-key", IssuedBy: catalogue.IssuedOutside},
|
||||
}
|
||||
|
||||
// holdingOf is one reading of the record without a store: the newest statements, the open conditions, the
|
||||
// catalogue, and each consumer's provider already looked up, as providerFor memoises it.
|
||||
// catalogue, what was given on the provider's machine (nothing), and each consumer's provider already looked up,
|
||||
// as providerFor memoises it.
|
||||
func holdingOf(open []conditions.Condition, healths map[string]inventory.NodeHealth, bound map[[3]string]catalogue.Chosen) *holding {
|
||||
h := &holding{ctx: context.Background(), healths: healths, open: open, providers: map[string]providerLookup{},
|
||||
shelf: map[string]catalogue.Manifest{"db": {Module: "db", Version: "1",
|
||||
Provides: []catalogue.Offer{{Name: "postgres-database", Scope: catalogue.ScopeMesh}}}}}
|
||||
shelf: map[string]catalogue.Manifest{"db": {Module: "db", Version: "1", OwnSecrets: dbSecrets,
|
||||
Provides: []catalogue.Offer{{Name: "postgres-database", Scope: catalogue.ScopeMesh}}}},
|
||||
waits: operatorWaitFacts{given: map[string]map[string]time.Time{"db@anchor": {}}}}
|
||||
for k, p := range bound {
|
||||
h.providers[k[0]+"\x00"+k[1]+"\x00"+k[2]] = providerLookup{p, true}
|
||||
}
|
||||
@@ -75,7 +84,7 @@ func TestAProviderWaitingForTheOperatorHoldsOnlyTheConsumersOfTheWaitingPart(t *
|
||||
healthyDB.State, healthyDB.Waits = link.StateHealthy, nil
|
||||
byCredential := holdingOf(nil, shopFailingBeside(waitingDatabaseFor(inventory.Wait{Part: "the server",
|
||||
Secret: "licence", What: "the licence key"})), shopOnTheDatabase)
|
||||
byCredential.shelf["db"] = catalogue.Manifest{Module: "db", Version: "1", Provides: []catalogue.Offer{{
|
||||
byCredential.shelf["db"] = catalogue.Manifest{Module: "db", Version: "1", OwnSecrets: dbSecrets, Provides: []catalogue.Offer{{
|
||||
Name: "postgres-database", Credential: &catalogue.OfferCredential{Own: "licence"}}}}
|
||||
for _, c := range []struct {
|
||||
name string
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
// Package beside is where a test finds another repository of the mesh it reads: the catalogue's manifests,
|
||||
// the node-engine's genesis template (novox/hq issue 432), the SDK's conformance fixtures (issue 449).
|
||||
// the node-engine's genesis template (novox/hq issue 432).
|
||||
//
|
||||
// A test used to read the checkout beside this one, `../../../mesh-catalog`, so its verdict depended on
|
||||
// whatever sat on the machine running it: a stale or dirty checkout failed it on a desktop, and where none
|
||||
@@ -10,8 +10,7 @@
|
||||
// test judges against those clones, so agreement with the other repository is checked where
|
||||
// `mesh/repo-check` runs; a repository missing there fails the test, never skips it. Which clones a
|
||||
// check gets is chosen from the inventory: mesh-catalog by the source of the `nats` module, mesh-host
|
||||
// by the source of `mesh-host`, and mesh-sdk as the controller's sibling at the commit the controller's
|
||||
// go.mod pins. In a mesh where either module has no source repository nothing is
|
||||
// by the source of `mesh-host`. In a mesh where either module has no source repository nothing is
|
||||
// cloned, and these tests fail loudly with "not beside this check": a cause in the setup, not in the
|
||||
// change. And in a delivery group that holds a mesh-catalog pull request, these tests read the
|
||||
// catalogue's main, not the group's head.
|
||||
|
||||
@@ -163,6 +163,21 @@ type Principal struct {
|
||||
// goes with the retired seat row.
|
||||
var seatsTheControllerAsks = []string{"node-build-agent", "mesh-build-machine"}
|
||||
|
||||
// TheControllersAsk says whether a message of stream, published on subject, is an ask the controller
|
||||
// itself makes: one on the accept subject of a seat in seatsTheControllerAsks, kept in that seat's own
|
||||
// work queue. Such an ask, given up on by the seat's worker, can be delivered again with the authority
|
||||
// the controller already holds — the publish on that seat's accepts and the stream API to remove the
|
||||
// original — and no other can (novox/hq issue 334, ADR 0264's consequences: no grant over the seats'
|
||||
// queues).
|
||||
func TheControllersAsk(stream, subject string) bool {
|
||||
for _, seat := range seatsTheControllerAsks {
|
||||
if stream == seatStreamName(seat) && strings.HasPrefix(subject, "mesh.seat."+seat+".accept.") {
|
||||
return true
|
||||
}
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
// SeatVerb is one verb of one seat, on every machine holding it.
|
||||
type SeatVerb struct{ Seat, Verb string }
|
||||
|
||||
|
||||
@@ -1,137 +0,0 @@
|
||||
package link
|
||||
|
||||
import (
|
||||
"os"
|
||||
"os/exec"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"testing"
|
||||
)
|
||||
|
||||
// Where the conformance fixtures are read from, held on a small SDK repository of the test's own: two
|
||||
// commits, the clone checked out at the older as the build seat leaves it at the running controller's pin
|
||||
// (novox/hq issue 449). The contents are the test's, because what is judged is which commit is read.
|
||||
type sdkBed struct {
|
||||
root, goMod, captured, capturedLog string
|
||||
older, newer string
|
||||
}
|
||||
|
||||
const fixtureName = "events/module-event.json"
|
||||
|
||||
func newSDKBed(t *testing.T) sdkBed {
|
||||
t.Helper()
|
||||
if _, err := exec.LookPath("git"); err != nil {
|
||||
t.Fatalf("git is needed to read the SDK at a pin, as a merge check does: %v", err)
|
||||
}
|
||||
b := sdkBed{root: t.TempDir()}
|
||||
clone := filepath.Join(b.root, "mesh-sdk")
|
||||
git := func(args ...string) string {
|
||||
t.Helper()
|
||||
cmd := exec.Command("git", append([]string{"-C", clone, "-c", "user.name=t", "-c", "user.email=t@t",
|
||||
"-c", "commit.gpgsign=false"}, args...)...)
|
||||
out, err := cmd.CombinedOutput()
|
||||
if err != nil {
|
||||
t.Fatalf("git %v: %v: %s", args, err, out)
|
||||
}
|
||||
return strings.TrimSpace(string(out))
|
||||
}
|
||||
write := func(dir, body string) {
|
||||
t.Helper()
|
||||
file := filepath.Join(dir, "conformance", filepath.FromSlash(fixtureName))
|
||||
if err := os.MkdirAll(filepath.Dir(file), 0o755); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := os.WriteFile(file, []byte(body), 0o644); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
if err := os.MkdirAll(clone, 0o755); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
git("init", "--quiet")
|
||||
write(clone, "older\n")
|
||||
git("add", "-A")
|
||||
git("commit", "--quiet", "-m", "older")
|
||||
b.older = git("rev-parse", "HEAD")
|
||||
write(clone, "newer\n")
|
||||
git("commit", "--quiet", "-am", "newer")
|
||||
b.newer = git("rev-parse", "HEAD")
|
||||
git("checkout", "--quiet", b.older)
|
||||
|
||||
other := t.TempDir()
|
||||
b.captured = filepath.Join(other, "mesh-sdk")
|
||||
write(b.captured, "older\n")
|
||||
b.capturedLog = filepath.Join(other, "CAPTURED")
|
||||
if err := os.WriteFile(b.capturedLog, []byte("mesh-sdk "+b.older+" conformance/events\n"), 0o644); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
b.goMod = filepath.Join(other, "go.mod")
|
||||
b.pin(t, b.newer)
|
||||
return b
|
||||
}
|
||||
|
||||
func (b sdkBed) pin(t *testing.T, commit string) {
|
||||
t.Helper()
|
||||
mod := "module x\n\nrequire (\n\t" + sdkModule + " v0.1.11-0.20261009143344-" + commit[:12] + "\n)\n"
|
||||
if err := os.WriteFile(b.goMod, []byte(mod), 0o644); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
|
||||
// A pull request moving the SDK is judged against the SDK it moves to: the clone is read at the tree's own
|
||||
// pin, not at what it is checked out at, the running controller's.
|
||||
func TestAMergeCheckReadsTheSDKAtTheTreesOwnPin(t *testing.T) {
|
||||
b := newSDKBed(t)
|
||||
raw, from, err := sdkFixture(b.root, b.goMod, b.captured, b.capturedLog, fixtureName)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if string(raw) != "newer\n" {
|
||||
t.Errorf("read %q from %s; the tree pins the newer commit", raw, from)
|
||||
}
|
||||
}
|
||||
|
||||
// A captured copy that is not the SDK's at the commit CAPTURED names fails a merge check.
|
||||
func TestAMergeCheckFailsACapturedCopyEditedByHand(t *testing.T) {
|
||||
b := newSDKBed(t)
|
||||
if err := os.WriteFile(filepath.Join(b.captured, "conformance", filepath.FromSlash(fixtureName)),
|
||||
[]byte("edited\n"), 0o644); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, _, err := sdkFixture(b.root, b.goMod, b.captured, b.capturedLog, fixtureName); err == nil ||
|
||||
!strings.Contains(err.Error(), "is not mesh-sdk's at") {
|
||||
t.Errorf("a hand-edited captured copy was accepted: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
// A pin the clone does not hold fails loudly and names it.
|
||||
func TestAMergeCheckFailsWithoutTheSDKAtItsPin(t *testing.T) {
|
||||
b := newSDKBed(t)
|
||||
b.pin(t, "0123456789ab0123456789ab0123456789ab0123")
|
||||
if _, _, err := sdkFixture(b.root, b.goMod, b.captured, b.capturedLog, fixtureName); err == nil ||
|
||||
!strings.Contains(err.Error(), "0123456789ab") {
|
||||
t.Errorf("a pin the clone lacks was not named: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
// Until mesh-controller #220 has rolled out the seat places no mesh-sdk beside a check: the captured copy is
|
||||
// read, and where it was read from says plainly that the clone did not judge it (novox/hq issue 449).
|
||||
func TestAMergeCheckWithoutTheCloneSaysItReadTheCapturedCopy(t *testing.T) {
|
||||
b := newSDKBed(t)
|
||||
raw, from, err := sdkFixture(t.TempDir(), b.goMod, b.captured, b.capturedLog, fixtureName)
|
||||
if err != nil || string(raw) != "older\n" {
|
||||
t.Fatalf("read %q, %v; the captured copy holds the older", raw, err)
|
||||
}
|
||||
if !strings.HasPrefix(from, "NOT JUDGED AGAINST THE CLONE") || !strings.Contains(from, b.older) {
|
||||
t.Errorf("the log does not say the clone was missing and which copy was read: %q", from)
|
||||
}
|
||||
}
|
||||
|
||||
// Away from a merge check the captured copy is read, and nothing else.
|
||||
func TestAwayFromACheckTheCapturedCopyIsRead(t *testing.T) {
|
||||
b := newSDKBed(t)
|
||||
raw, _, err := sdkFixture("", b.goMod, b.captured, b.capturedLog, fixtureName)
|
||||
if err != nil || string(raw) != "older\n" {
|
||||
t.Errorf("read %q, %v; the captured copy holds the older", raw, err)
|
||||
}
|
||||
}
|
||||
@@ -1,40 +1,19 @@
|
||||
package link
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"os"
|
||||
"os/exec"
|
||||
"path"
|
||||
"path/filepath"
|
||||
"regexp"
|
||||
"strings"
|
||||
"sync"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/novox/mesh-controller/internal/beside"
|
||||
)
|
||||
|
||||
// The Go implementation, held to the shared fixtures (novox/hq ADR 0074, design 19).
|
||||
//
|
||||
// **Read from the SDK's own conformance directory, never copied by hand**: a fixture written into each
|
||||
// implementation is two fixtures, and two fixtures drift, which is the failure the suite exists to prevent.
|
||||
// Which SDK is read is the one this tree's go.mod pins, never the one a desktop happens to hold beside this
|
||||
// checkout (novox/hq issue 449):
|
||||
//
|
||||
// - in a merge check, the fixture as it stands at that pin in the mesh-sdk clone the build seat places
|
||||
// beside the check — read with `git show`, because the clone is checked out at the *running*
|
||||
// controller's pin, and a pull request moving the SDK must be judged against the SDK it moves to, or it
|
||||
// passes the check and fails everywhere after it rolls out. A pin the clone lacks fails the test; a
|
||||
// missing clone reads the captured copy and says so loudly, until mesh-controller #220 has rolled out
|
||||
// and the build seat clones mesh-sdk (then its follow-up makes a missing clone fail again);
|
||||
// - elsewhere, the copy captured in testdata/beside at the commit testdata/beside/CAPTURED names, held to
|
||||
// go.mod by TestTheCapturedSDKIsTheOneGoModPins.
|
||||
//
|
||||
// In a merge check the captured copy is also compared with the clone at the captured commit, byte for byte,
|
||||
// so a hand-edited copy cannot pass for the SDK's.
|
||||
// **Read from the sdk's conformance directory by sibling path**, the way the lab finds its
|
||||
// siblings — deliberately not copied here. A fixture copied into each implementation is two
|
||||
// fixtures, and two fixtures drift, which is the exact failure the suite exists to prevent.
|
||||
type fixture struct {
|
||||
Name string `json:"name"`
|
||||
Given struct {
|
||||
@@ -51,27 +30,12 @@ type fixture struct {
|
||||
} `json:"wire"`
|
||||
}
|
||||
|
||||
// sdkModule is the Go module of the SDK the controller is built against.
|
||||
const sdkModule = "git.novox.be/novox/mesh-sdk/go"
|
||||
|
||||
// The tree's go.mod, and the captured copy, as this package finds them.
|
||||
var (
|
||||
goModFile = filepath.Join("..", "..", "go.mod")
|
||||
capturedSDK = filepath.Join(beside.Captured(), "mesh-sdk")
|
||||
capturedLog = filepath.Join(beside.Captured(), "CAPTURED")
|
||||
)
|
||||
|
||||
func loadFixture(t *testing.T, name string) fixture {
|
||||
t.Helper()
|
||||
raw, from, err := sdkFixture(os.Getenv(beside.Env), goModFile, capturedSDK, capturedLog, name)
|
||||
path := filepath.Join("..", "..", "..", "mesh-sdk", "conformance", name)
|
||||
raw, err := os.ReadFile(path)
|
||||
if err != nil {
|
||||
// Failed, never skipped: a skip here passed the suite with nothing judged.
|
||||
t.Fatal(err)
|
||||
}
|
||||
t.Logf("%s: %s", name, from)
|
||||
if strings.HasPrefix(from, notJudged) {
|
||||
// Said on the check's own output too, where `go test` without -v prints no log of a passing test.
|
||||
saidNotJudged.Do(func() { fmt.Fprintln(os.Stderr, "internal/link conformance: "+from) })
|
||||
t.Skipf("the sdk's conformance fixtures are not beside this checkout: %v", err)
|
||||
}
|
||||
var f fixture
|
||||
if err := json.Unmarshal(raw, &f); err != nil {
|
||||
@@ -80,141 +44,6 @@ func loadFixture(t *testing.T, name string) fixture {
|
||||
return f
|
||||
}
|
||||
|
||||
// notJudged opens what a merge check without the SDK's clone says; saidNotJudged says it once on stderr.
|
||||
const notJudged = "NOT JUDGED AGAINST THE CLONE"
|
||||
|
||||
var saidNotJudged sync.Once
|
||||
|
||||
// sdkFixture is one conformance fixture and where it was read: from the mesh-sdk clone in root (a merge
|
||||
// check's MESH_CHECK_BESIDE) at the pin of goMod, or, with no root, from the captured copy.
|
||||
func sdkFixture(root, goMod, captured, capturedLog, name string) ([]byte, string, error) {
|
||||
file := path.Join("conformance", name)
|
||||
readCaptured := func() ([]byte, error) {
|
||||
raw, err := os.ReadFile(filepath.Join(captured, filepath.FromSlash(file)))
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("the captured SDK's %s: %w", file, err)
|
||||
}
|
||||
return raw, nil
|
||||
}
|
||||
if root == "" {
|
||||
raw, err := readCaptured()
|
||||
return raw, "the copy captured in testdata/beside (see CAPTURED); a merge check reads the SDK's clone " +
|
||||
"at this tree's pin", err
|
||||
}
|
||||
clone := filepath.Join(root, "mesh-sdk")
|
||||
if _, err := os.Stat(clone); err != nil {
|
||||
// A build seat asked by a controller from before mesh-controller #220 places no mesh-sdk beside the
|
||||
// check, so #220 could never pass its own check if this failed. Until #220 has rolled out, the
|
||||
// captured copy is read and the log says so loudly; the follow-up of novox/hq issue 449 makes this
|
||||
// a failure again.
|
||||
at, cerr := capturedCommit(capturedLog)
|
||||
if cerr != nil {
|
||||
return nil, "", cerr
|
||||
}
|
||||
raw, rerr := readCaptured()
|
||||
return raw, fmt.Sprintf(notJudged+": the seat placed no mesh-sdk beside this check in "+
|
||||
"%s=%s (it does once mesh-controller #220 has rolled out); read the captured copy at %s", beside.Env,
|
||||
root, at), rerr
|
||||
}
|
||||
ref, err := sdkRef(goMod)
|
||||
if err != nil {
|
||||
return nil, "", err
|
||||
}
|
||||
raw, err := gitShow(clone, ref, file)
|
||||
if err != nil {
|
||||
return nil, "", fmt.Errorf("the mesh-sdk clone beside this check has no %s at %s, the SDK this tree's "+
|
||||
"go.mod pins: %w", file, ref, err)
|
||||
}
|
||||
// The captured copy is the SDK's own, never edited by hand: the clone at the captured commit says so.
|
||||
at, err := capturedCommit(capturedLog)
|
||||
if err != nil {
|
||||
return nil, "", err
|
||||
}
|
||||
theirs, err := gitShow(clone, at, file)
|
||||
if err != nil {
|
||||
return nil, "", fmt.Errorf("the mesh-sdk clone has no %s at %s, the commit CAPTURED names: %w", file, at, err)
|
||||
}
|
||||
ours, err := os.ReadFile(filepath.Join(captured, filepath.FromSlash(file)))
|
||||
if err != nil {
|
||||
return nil, "", fmt.Errorf("the captured SDK's %s: %w", file, err)
|
||||
}
|
||||
if !bytes.Equal(ours, theirs) {
|
||||
return nil, "", fmt.Errorf("the captured %s is not mesh-sdk's at %s, the commit CAPTURED names: it was "+
|
||||
"edited or captured wrong; capture it again as CAPTURED says", file, at)
|
||||
}
|
||||
return raw, "mesh-sdk cloned beside this check, at " + ref + " (this tree's go.mod)", nil
|
||||
}
|
||||
|
||||
// gitShow is one file of a repository at a ref. safe.directory, because the clone is the build seat's and
|
||||
// the test may run as another user.
|
||||
func gitShow(repository, ref, file string) ([]byte, error) {
|
||||
cmd := exec.Command("git", "-c", "safe.directory=*", "-C", repository, "show", ref+":"+file)
|
||||
var stderr bytes.Buffer
|
||||
cmd.Stderr = &stderr
|
||||
out, err := cmd.Output()
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("git show %s:%s: %v: %s", ref, file, err, strings.TrimSpace(stderr.String()))
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
|
||||
// pseudoCommit is the commit a Go pseudo-version names: v0.1.11-0.20261009143344-f047d0a4a970 → f047d0a4a970.
|
||||
var pseudoCommit = regexp.MustCompile(`-([0-9a-f]{12})$`)
|
||||
|
||||
// sdkRef is the SDK repository's ref a go.mod pins: a pseudo-version's commit, or a release's tag (the SDK
|
||||
// tags its Go module under go/).
|
||||
func sdkRef(goMod string) (string, error) {
|
||||
raw, err := os.ReadFile(goMod)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
version := ""
|
||||
for _, line := range strings.Split(string(raw), "\n") {
|
||||
fields := strings.Fields(strings.TrimPrefix(strings.TrimSpace(line), "require "))
|
||||
if len(fields) >= 2 && fields[0] == sdkModule {
|
||||
version = fields[1]
|
||||
}
|
||||
}
|
||||
if m := pseudoCommit.FindStringSubmatch(version); m != nil {
|
||||
return m[1], nil
|
||||
}
|
||||
if strings.HasPrefix(version, "v") {
|
||||
return "go/" + version, nil
|
||||
}
|
||||
return "", fmt.Errorf("%s pins no version of %s that names a commit or a tag", goMod, sdkModule)
|
||||
}
|
||||
|
||||
// capturedCommit is the commit testdata/beside/CAPTURED names for mesh-sdk.
|
||||
func capturedCommit(capturedLog string) (string, error) {
|
||||
raw, err := os.ReadFile(capturedLog)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
at := regexp.MustCompile(`(?m)^mesh-sdk\s+([0-9a-f]{40})\s`).FindSubmatch(raw)
|
||||
if at == nil {
|
||||
return "", fmt.Errorf("%s names no commit for mesh-sdk", capturedLog)
|
||||
}
|
||||
return string(at[1]), nil
|
||||
}
|
||||
|
||||
// The captured SDK is the one this controller is built against: when go.mod moves the SDK, the copy moves
|
||||
// with it, or the tests away from a merge check judge an SDK the controller no longer uses (novox/hq issue
|
||||
// 449).
|
||||
func TestTheCapturedSDKIsTheOneGoModPins(t *testing.T) {
|
||||
pinned, err := sdkRef(goModFile)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
at, err := capturedCommit(capturedLog)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if strings.HasPrefix(pinned, "go/") || !strings.HasPrefix(at, pinned) {
|
||||
t.Errorf("the SDK is captured at %s but go.mod pins %s: capture it again at the pinned commit, as "+
|
||||
"testdata/beside/CAPTURED says", at, pinned)
|
||||
}
|
||||
}
|
||||
|
||||
// Every header the fixture requires is one this implementation actually sets.
|
||||
func TestTheGoEmitterSetsEveryRequiredHeader(t *testing.T) {
|
||||
f := loadFixture(t, "events/module-event.json")
|
||||
|
||||
@@ -62,6 +62,9 @@ type DeadLetter struct {
|
||||
// taken. The record of it is kept all the same, so it is said and dropped, never silently missing.
|
||||
Lost string `json:"lost,omitempty"`
|
||||
Size int `json:"size"`
|
||||
// Original says what became of the ask it was kept from, in its seat's work queue, when it was
|
||||
// delivered again (novox/hq issue 334).
|
||||
Original string `json:"original,omitempty"`
|
||||
// Body and Headers are the message itself, given only for one dead letter asked by its id; a header
|
||||
// with several values keeps them all.
|
||||
Body string `json:"body,omitempty"`
|
||||
@@ -250,9 +253,15 @@ func DeadLetterNamed(js nats.JetStreamContext, id uint64) (DeadLetter, error) {
|
||||
|
||||
// AgainTo is where a kept message is delivered again so that only the consumer that gave it up gets
|
||||
// it: an event under that consumer's own again subject on EVENTS — and only when the consumer exists and
|
||||
// filters that subject, so a message is never let go as delivered while nobody receives it. Any other
|
||||
// stream's message is refused, with why: an ask given up on by a seat's worker is not delivered again yet
|
||||
// (novox/hq issue 330's follow-up), since publishing it again leaves the original stuck in the queue.
|
||||
// filters that subject, so a message is never let go as delivered while nobody receives it.
|
||||
//
|
||||
// An ask the controller itself made (broker.TheControllersAsk) goes back on its own subject: a seat's
|
||||
// work queue has one worker per subject, so only the worker that gave it up takes it, and DeliverAgain
|
||||
// removes the original from the queue first, so there are never two (novox/hq issue 334); only while the
|
||||
// worker that gave it up is on the bus. Any other stream's
|
||||
// message is refused, with why: an ask to a seat the controller does not ask would need a publish it is
|
||||
// not granted, in its asker's name (ADR 0264's consequences, ADR 0259 §3), and a stream that is neither
|
||||
// would reach every consumer of its subject.
|
||||
func AgainTo(js nats.JetStreamContext, d DeadLetter) (string, error) {
|
||||
switch {
|
||||
case d.Lost != "":
|
||||
@@ -260,10 +269,22 @@ func AgainTo(js nats.JetStreamContext, d DeadLetter) (string, error) {
|
||||
case d.Subject == "":
|
||||
return "", fmt.Errorf("dead letter %d does not say the subject it was published on, so it cannot be "+
|
||||
"delivered again. Drop it", d.ID)
|
||||
case broker.TheControllersAsk(d.Stream, d.Subject):
|
||||
// The seat's worker takes it, and only while it is on the bus: without one the queue would keep
|
||||
// the ask for a holder that may never come, and it would be let go as delivered meanwhile.
|
||||
if _, err := js.ConsumerInfo(d.Stream, d.Consumer); errors.Is(err, nats.ErrConsumerNotFound) {
|
||||
return "", fmt.Errorf("dead letter %d was given up by %s, which is no longer on the bus, so the ask "+
|
||||
"would only wait in %s for a holder. Nothing was done, and it is still kept: deliver it again once "+
|
||||
"the seat has a holder, or drop it", d.ID, d.Who, d.Stream)
|
||||
} else if err != nil {
|
||||
return "", fmt.Errorf("whether %s can receive dead letter %d cannot be read: %w", d.Who, d.ID, err)
|
||||
}
|
||||
return d.Subject, nil
|
||||
case d.Stream != broker.EventsStream:
|
||||
return "", fmt.Errorf("dead letter %d is from %s, and only an event is delivered again: publishing it "+
|
||||
"again would reach every consumer of its subject, or leave the original in its queue. Drop it, and "+
|
||||
"have its sender say it again", d.ID, d.Stream)
|
||||
return "", fmt.Errorf("dead letter %d is from %s, and only an event or an ask the controller made is "+
|
||||
"delivered again: publishing it again would reach every consumer of its subject, or need a grant the "+
|
||||
"controller does not hold to speak for its asker. Drop it, and have its sender say it again",
|
||||
d.ID, d.Stream)
|
||||
}
|
||||
info, err := js.ConsumerInfo(d.Stream, d.Consumer)
|
||||
if errors.Is(err, nats.ErrConsumerNotFound) {
|
||||
@@ -284,7 +305,7 @@ func AgainTo(js nats.JetStreamContext, d DeadLetter) (string, error) {
|
||||
}
|
||||
|
||||
// DeliverAgain hands a kept message to the consumer that gave it up, and nobody else, then removes it
|
||||
// from DEAD_LETTERS. The message carries its own headers and AgainHeader; its de-duplication id is the
|
||||
// from DEAD_LETTERS; an ask's original is removed from its work queue before. The message carries its own headers and AgainHeader; its de-duplication id is the
|
||||
// kept copy's, so asking twice delivers it once.
|
||||
func DeliverAgain(js nats.JetStreamContext, id uint64) (DeadLetter, string, error) {
|
||||
d, err := DeadLetterNamed(js, id)
|
||||
@@ -295,6 +316,9 @@ func DeliverAgain(js nats.JetStreamContext, id uint64) (DeadLetter, string, erro
|
||||
if err != nil {
|
||||
return d, "", err
|
||||
}
|
||||
if d.Stream != broker.EventsStream {
|
||||
d.Original = removeOriginal(js, d)
|
||||
}
|
||||
again := &nats.Msg{Subject: to, Header: nats.Header{}, Data: []byte(d.Body)}
|
||||
for k, v := range d.Headers {
|
||||
again.Header[k] = append([]string(nil), v...)
|
||||
@@ -302,6 +326,10 @@ func DeliverAgain(js nats.JetStreamContext, id uint64) (DeadLetter, string, erro
|
||||
again.Header.Set(AgainHeader, strconv.FormatUint(id, 10))
|
||||
again.Header.Set(nats.MsgIdHdr, "again."+broker.DeadLettersStream+"."+strconv.FormatUint(id, 10))
|
||||
if _, err := js.PublishMsg(again); err != nil {
|
||||
if d.Original != "" {
|
||||
return d, to, fmt.Errorf("dead letter %d could not be delivered again on %s: %w; it is still kept, so "+
|
||||
"delivering it again tries once more. Its original: %s", id, to, err, d.Original)
|
||||
}
|
||||
return d, to, fmt.Errorf("dead letter %d could not be delivered again on %s: %w; it is still kept", id, to, err)
|
||||
}
|
||||
if err := js.DeleteMsg(broker.DeadLettersStream, id); err != nil && !errors.Is(err, nats.ErrMsgNotFound) {
|
||||
@@ -311,6 +339,34 @@ func DeliverAgain(js nats.JetStreamContext, id uint64) (DeadLetter, string, erro
|
||||
return d, to, nil
|
||||
}
|
||||
|
||||
// removeOriginal takes the ask a dead letter was kept from out of its seat's work queue, and says what
|
||||
// became of it. Given up on, the original is never acknowledged and would stay beside its copy until the
|
||||
// stream's age drops it (novox/hq issue 334). It is removed before the copy is published, so a copy that
|
||||
// cannot be published leaves the kept one to try again, and never two.
|
||||
//
|
||||
// **Only when the queue still holds that same message**: its subject and the time it was stored are the
|
||||
// dead letter's. A seat's stream deleted and made again — the build handover deletes one (builds.go) —
|
||||
// numbers from one again, and an old dead letter's sequence may then name another, live ask; deleting by
|
||||
// the number alone would drop that one silently. Anything else is said, never an error: the original is
|
||||
// gone or is not this one, and the copy is the only one there will be.
|
||||
func removeOriginal(js nats.JetStreamContext, d DeadLetter) string {
|
||||
held, err := js.GetMsg(d.Stream, d.Sequence)
|
||||
switch {
|
||||
case errors.Is(err, nats.ErrMsgNotFound):
|
||||
return fmt.Sprintf("%s no longer held message %d, so there was nothing to remove", d.Stream, d.Sequence)
|
||||
case err != nil:
|
||||
return fmt.Sprintf("message %d of %s could not be read, so it was left as it is: %v", d.Sequence, d.Stream, err)
|
||||
case d.Published.IsZero() || held.Subject != d.Subject || !held.Time.Equal(d.Published):
|
||||
return fmt.Sprintf("message %d of %s is another message now (%s, stored %s), so it was left as it is; "+
|
||||
"the one given up on is gone", d.Sequence, d.Stream, held.Subject, held.Time.UTC().Format(time.RFC3339))
|
||||
}
|
||||
if err := js.DeleteMsg(d.Stream, d.Sequence); err != nil && !errors.Is(err, nats.ErrMsgNotFound) {
|
||||
return fmt.Sprintf("message %d of %s could not be removed, so it stays beside its copy until the "+
|
||||
"stream's age drops it; nothing delivers it again: %v", d.Sequence, d.Stream, err)
|
||||
}
|
||||
return fmt.Sprintf("message %d of %s, the one given up on, was removed from the queue", d.Sequence, d.Stream)
|
||||
}
|
||||
|
||||
// DropDeadLetter removes a kept message for good.
|
||||
func DropDeadLetter(js nats.JetStreamContext, id uint64) (DeadLetter, error) {
|
||||
d, err := DeadLetterNamed(js, id)
|
||||
|
||||
@@ -0,0 +1,217 @@
|
||||
package link
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
"github.com/nats-io/nats.go"
|
||||
"github.com/nats-io/nats.go/jetstream"
|
||||
|
||||
"github.com/novox/mesh-controller/internal/broker"
|
||||
)
|
||||
|
||||
// What a seat's worker gave up on, when the ask is one the controller itself makes (novox/hq issue 334).
|
||||
//
|
||||
// The controller already publishes on the accept subjects of the seats it asks (broker's
|
||||
// seatsTheControllerAsks) and reaches the stream API, so delivering its own ask again needs no grant it
|
||||
// does not hold: the original is removed from the work queue by its sequence, and the kept copy is
|
||||
// published on the ask's own subject, where the seat's one worker takes it. An ask to any other seat is
|
||||
// still refused, and stays kept: delivering it would need a publish the controller is not granted
|
||||
// (ADR 0264's consequences), in the asker's name (ADR 0259 §3).
|
||||
|
||||
// theBuildWorker is the build seat's worker as a holder pulls from it.
|
||||
func theBuildWorker(t *testing.T, js *broker.JetStream) jetstream.Consumer {
|
||||
t.Helper()
|
||||
api, err := jetstream.New(js.Conn())
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
stream := "SEAT_NODE_BUILD_AGENT"
|
||||
worker, err := api.Consumer(t.Context(), stream, stream+"_worker")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
return worker
|
||||
}
|
||||
|
||||
func TestAnAskTheControllerMadeIsDeliveredAgainToTheSeatsWorker(t *testing.T) {
|
||||
js := aBusWithTheBuildRole(t)
|
||||
keeping(t, js)
|
||||
worker := theBuildWorker(t, js)
|
||||
|
||||
subject := BuildWorkOf(TheBuildMachine)
|
||||
ask := &nats.Msg{Subject: subject, Data: []byte(`{"module":"x"}`), Header: nats.Header{}}
|
||||
ask.Header.Set(nats.MsgIdHdr, "build-1")
|
||||
if _, err := js.Context().PublishMsg(ask); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
d := givenUp(t, js, worker)
|
||||
if d.Stream != "SEAT_NODE_BUILD_AGENT" || d.Subject != subject || d.Lost != "" {
|
||||
t.Fatalf("kept as %+v", d)
|
||||
}
|
||||
if _, err := js.Context().GetMsg(d.Stream, d.Sequence); err != nil {
|
||||
t.Fatalf("the original is not in the work queue before it is delivered again: %v", err)
|
||||
}
|
||||
|
||||
_, to, err := DeliverAgain(js.Context(), d.ID)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if to != subject {
|
||||
t.Fatalf("delivered again on %s, not the ask's own subject %s", to, subject)
|
||||
}
|
||||
if _, err := js.Context().GetMsg(d.Stream, d.Sequence); !errors.Is(err, nats.ErrMsgNotFound) {
|
||||
t.Fatalf("the original is still in the work queue beside its copy: %v", err)
|
||||
}
|
||||
again := next(t, worker)
|
||||
if again == nil {
|
||||
t.Fatal("the seat's worker was not handed the ask again")
|
||||
}
|
||||
if string(again.Data()) != `{"module":"x"}` || again.Headers().Get(AgainHeader) == "" {
|
||||
t.Fatalf("handed again as %s %v", again.Data(), again.Headers())
|
||||
}
|
||||
_ = again.Ack()
|
||||
if m := next(t, worker); m != nil {
|
||||
t.Fatalf("the worker was handed it twice: %s", m.Subject())
|
||||
}
|
||||
if _, err := DeadLetterNamed(js.Context(), d.ID); !errors.Is(err, ErrNoDeadLetter) {
|
||||
t.Fatalf("still kept after it was delivered again: %v", err)
|
||||
}
|
||||
if _, _, err := DeliverAgain(js.Context(), d.ID); !errors.Is(err, ErrNoDeadLetter) {
|
||||
t.Fatalf("a second delivery answered %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
// An ask to a seat the controller does not ask is refused, says why, and nothing is done.
|
||||
func TestAnAskTheControllerDidNotMakeIsStillOnlyDropped(t *testing.T) {
|
||||
js := aBus(t)
|
||||
for _, d := range []DeadLetter{
|
||||
{ID: 2, Stream: "SEAT_TELEGRAM_SENDER", Consumer: "SEAT_TELEGRAM_SENDER_worker",
|
||||
Subject: "mesh.seat.telegram-sender.accept.send"},
|
||||
// The subject of a seat the controller asks, on a stream that is not that seat's queue.
|
||||
{ID: 3, Stream: "SEAT_TELEGRAM_SENDER", Consumer: "SEAT_TELEGRAM_SENDER_worker",
|
||||
Subject: "mesh.seat.node-build-agent.accept.build"},
|
||||
} {
|
||||
if to, err := AgainTo(js.Context(), d); err == nil {
|
||||
t.Errorf("%s on %s was given %s to be delivered again on", d.Subject, d.Stream, to)
|
||||
} else if !strings.Contains(err.Error(), "Drop it") {
|
||||
t.Errorf("refused without saying what to do: %v", err)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// aBuildAskGivenUp publishes one build ask and lets the worker give it up.
|
||||
func aBuildAskGivenUp(t *testing.T, js *broker.JetStream, worker jetstream.Consumer, body string) DeadLetter {
|
||||
t.Helper()
|
||||
if _, err := js.Context().Publish(BuildWorkOf(TheBuildMachine), []byte(body)); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
return givenUp(t, js, worker)
|
||||
}
|
||||
|
||||
// A copy that cannot be published: the original is already out of the queue, the dead letter is kept and
|
||||
// the answer says both; asked again once it can be published, it is delivered once.
|
||||
func TestAnAskWhoseCopyIsRefusedStaysKeptAndIsDeliveredOnTheNextTry(t *testing.T) {
|
||||
js := aBusWithTheBuildRole(t)
|
||||
keeping(t, js)
|
||||
worker := theBuildWorker(t, js)
|
||||
d := aBuildAskGivenUp(t, js, worker, `{"module":"x"}`)
|
||||
|
||||
// The queue stops taking the ask's subject, so the copy's publish is refused by the server.
|
||||
info, err := js.Context().StreamInfo(d.Stream)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
cfg := info.Config
|
||||
taking := cfg.Subjects
|
||||
cfg.Subjects = []string{"mesh.seat." + TheBuildMachine + ".accept.nothing"}
|
||||
if _, err := js.Context().UpdateStream(&cfg); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
_, _, err = DeliverAgain(js.Context(), d.ID)
|
||||
if err == nil || !strings.Contains(err.Error(), "still kept") || !strings.Contains(err.Error(), "was removed") {
|
||||
t.Fatalf("a refused copy answered %v", err)
|
||||
}
|
||||
if _, err := js.Context().GetMsg(d.Stream, d.Sequence); !errors.Is(err, nats.ErrMsgNotFound) {
|
||||
t.Fatalf("the original is still in the queue: %v", err)
|
||||
}
|
||||
if _, err := DeadLetterNamed(js.Context(), d.ID); err != nil {
|
||||
t.Fatalf("a refused copy let the dead letter go: %v", err)
|
||||
}
|
||||
|
||||
cfg.Subjects = taking
|
||||
if _, err := js.Context().UpdateStream(&cfg); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
retried, _, err := DeliverAgain(js.Context(), d.ID)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if !strings.Contains(retried.Original, "no longer held") {
|
||||
t.Fatalf("the retry said of the original: %q", retried.Original)
|
||||
}
|
||||
if again := next(t, worker); again == nil || string(again.Data()) != `{"module":"x"}` {
|
||||
t.Fatal("the retry did not hand the ask to the worker")
|
||||
} else {
|
||||
_ = again.Ack()
|
||||
}
|
||||
if m := next(t, worker); m != nil {
|
||||
t.Fatalf("handed twice: %s", m.Data())
|
||||
}
|
||||
}
|
||||
|
||||
// A queue made again numbers from one: the old dead letter's sequence then names a live ask, which is left
|
||||
// alone.
|
||||
func TestAnAskWhoseSequenceNamesAnotherMessageLeavesThatOneAlone(t *testing.T) {
|
||||
js := aBusWithTheBuildRole(t)
|
||||
keeping(t, js)
|
||||
d := aBuildAskGivenUp(t, js, theBuildWorker(t, js), `{"module":"old"}`)
|
||||
|
||||
// The build handover's way: the seat's queue deleted and made again, with its worker.
|
||||
if err := js.Context().DeleteStream(d.Stream); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := broker.RaiseSeats(js, []broker.DeclaredSeat{{Name: TheBuildMachine, Accepts: []string{"build"},
|
||||
Emits: []string{"built"}}}, map[string]broker.Holder{TheBuildMachine: {Node: "anchor", Module: "builder"}}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
live, err := js.Context().Publish(BuildWorkOf(TheBuildMachine), []byte(`{"module":"live"}`))
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if live.Sequence != d.Sequence {
|
||||
t.Fatalf("the live ask is message %d, the dead letter names %d: the test does not set up the collision",
|
||||
live.Sequence, d.Sequence)
|
||||
}
|
||||
|
||||
delivered, _, err := DeliverAgain(js.Context(), d.ID)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if !strings.Contains(delivered.Original, "another message") {
|
||||
t.Fatalf("the answer said of the original: %q", delivered.Original)
|
||||
}
|
||||
if held, err := js.Context().GetMsg(d.Stream, d.Sequence); err != nil || string(held.Data) != `{"module":"live"}` {
|
||||
t.Fatalf("the live ask at that sequence was touched: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
// No worker on the seat's queue: refused, kept, and nothing published.
|
||||
func TestAnAskIsNotDeliveredAgainWhileTheSeatHasNoWorker(t *testing.T) {
|
||||
js := aBusWithTheBuildRole(t)
|
||||
keeping(t, js)
|
||||
d := aBuildAskGivenUp(t, js, theBuildWorker(t, js), `{"module":"x"}`)
|
||||
if err := js.Context().DeleteConsumer(d.Stream, d.Consumer); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, _, err := DeliverAgain(js.Context(), d.ID); err == nil || !strings.Contains(err.Error(), "wait in") {
|
||||
t.Fatalf("delivered with no worker: %v", err)
|
||||
}
|
||||
if _, err := js.Context().GetMsg(d.Stream, d.Sequence); err != nil {
|
||||
t.Fatalf("a refusal removed the original: %v", err)
|
||||
}
|
||||
if _, err := DeadLetterNamed(js.Context(), d.ID); err != nil {
|
||||
t.Fatalf("a refusal let the dead letter go: %v", err)
|
||||
}
|
||||
}
|
||||
Vendored
+1
-6
@@ -9,7 +9,6 @@ a second place to clean besides the catalogue itself.
|
||||
|
||||
mesh-catalog b9de001833b1b61b297182b8e3bcdae1cbdeace9 modules/*/module.json, modules/nats/Dockerfile
|
||||
mesh-host bd5cc6980419c1bd4824be6d58dfebd2381a9de3 examples/foundation-first-node-nats.lock
|
||||
mesh-sdk f047d0a4a9702f5d7f3d7eadd3e62496311acc00 conformance/events/module-event.json
|
||||
|
||||
To move them, from this repository's root, with the two repositories checked out beside it:
|
||||
|
||||
@@ -20,8 +19,4 @@ To move them, from this repository's root, with the two repositories checked out
|
||||
find testdata/beside -name module.json -exec mv {} {}.captured \;
|
||||
git -C ../mesh-host archive <commit> examples/foundation-first-node-nats.lock | tar -x -C testdata/beside/mesh-host
|
||||
|
||||
and write the commits here. The SDK is captured at the commit go.mod pins for git.novox.be/novox/mesh-sdk/go
|
||||
(the last part of its pseudo-version; novox/hq issue 449), which internal/link's conformance test holds it to:
|
||||
|
||||
rm -rf testdata/beside/mesh-sdk && mkdir -p testdata/beside/mesh-sdk
|
||||
git -C ../mesh-sdk archive <commit> conformance/events | tar -x -C testdata/beside/mesh-sdk
|
||||
and write the commits here.
|
||||
|
||||
@@ -1,44 +0,0 @@
|
||||
{
|
||||
"capability": "events",
|
||||
"name": "a module emits an event",
|
||||
"why": "The envelope is what two implementations can disagree about without either failing: a missing header, a header spelled differently, or a body nested where metadata belongs. None of those stop a mesh running; they stop it reacting.",
|
||||
"given": {
|
||||
"module": "shop",
|
||||
"node": "one",
|
||||
"key": "order.placed",
|
||||
"body": { "id": "a1", "total": 12 },
|
||||
"headers": {
|
||||
"x-event-id": "0123456789abcdef0123456789abcdef",
|
||||
"x-source": "shop",
|
||||
"x-node": "one",
|
||||
"x-time": "2026-09-26T12:00:00Z",
|
||||
"content-type": "application/json"
|
||||
}
|
||||
},
|
||||
"wire": {
|
||||
"subject": "mesh.mod.shop.event.order.placed",
|
||||
"requiredHeaders": ["x-event-id", "x-source", "x-node", "x-time", "content-type"],
|
||||
"optionalHeaders": ["x-causation-id", "x-schema"],
|
||||
"headerFormats": {
|
||||
"x-time": "RFC3339",
|
||||
"content-type": "application/json",
|
||||
"x-event-id": "^[0-9a-f]{32}$"
|
||||
},
|
||||
"payloadIs": "the body alone, not the envelope",
|
||||
"keyRecoveredFrom": "the subject, after the .event. token"
|
||||
},
|
||||
"refuses": [
|
||||
{
|
||||
"what": "an envelope nested in the payload",
|
||||
"why": "an implementation that publishes the whole envelope as the body passes all of its own tests and is unreadable to every other"
|
||||
},
|
||||
{
|
||||
"what": "a missing x-event-id",
|
||||
"why": "delivery is at-least-once and only the emitter can say which of two messages is a redelivery"
|
||||
},
|
||||
{
|
||||
"what": "an x-source that differs from the subject's module",
|
||||
"why": "the bus enforces the namespace, so a disagreement means the envelope is lying about its origin"
|
||||
}
|
||||
]
|
||||
}
|
||||
Reference in New Issue
Block a user