From 03aa8413c28721d3217d0cfc97040591244a1fa9 Mon Sep 17 00:00:00 2001 From: Jason Stiebs Date: Fri, 11 Sep 2026 12:24:40 -0500 Subject: [PATCH 1/8] Require completed checker evidence for Jepsen qualification --- test/jepsen/qualification-result-test.sh | 36 ++++++++++++++++++ test/jepsen/qualification-result.sh | 19 ++++++++++ test/jepsen/qualify.sh | 14 +++---- test/jepsen/src/group/jepsen/core.clj | 3 +- .../jepsen/src/group/jepsen/qualification.clj | 37 +++++++++++++++++++ .../test/group/jepsen/qualification_test.clj | 26 +++++++++++++ 6 files changed, 126 insertions(+), 9 deletions(-) create mode 100755 test/jepsen/qualification-result-test.sh create mode 100644 test/jepsen/qualification-result.sh create mode 100644 test/jepsen/src/group/jepsen/qualification.clj create mode 100644 test/jepsen/test/group/jepsen/qualification_test.clj diff --git a/test/jepsen/qualification-result-test.sh b/test/jepsen/qualification-result-test.sh new file mode 100755 index 0000000..3dda51b --- /dev/null +++ b/test/jepsen/qualification-result-test.sh @@ -0,0 +1,36 @@ +#!/usr/bin/env bash +set -euo pipefail +script_dir="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" +source "${script_dir}/qualification-result.sh" +work="$(mktemp -d "${script_dir}/qualification-test.XXXXXX")" +trap 'rm -rf "${work}"' EXIT + +probe() { + local want="$1" expectation="$2" mode="$3" status="$4" valid="$5" qualified="$6" + local result="${work}/result" + rm -f "${result}" + # A bounded stand-in for the CLI: no Docker, JVM or network required. + if [[ "${valid}" != absent ]]; then + printf 'group-qualification-v1\t%s\t%s\t%s\n' "${mode}" "${valid}" "${qualified}" >"${result}" + fi + local actual=reject + if qualification_result "${expectation}" "${mode}" "${status}" "${result}"; then + actual=accept + fi + [[ "${actual}" == "${want}" ]] || { echo "unexpected classification"; exit 1; } +} + +probe accept pass none 0 true true +for mode in unexpected-death internal-index cursor-marker registry-projection terminal-unavailable; do + probe accept fail "${mode}" 1 false true + probe reject fail "${mode}" 1 false false +done +probe reject pass none 1 false false +probe reject pass none 1 absent false +probe reject fail internal-index 0 true true +for status in 1 124 125 137 127; do + probe reject fail internal-index "${status}" absent false +done +probe reject fail internal-index 124 false true +probe reject fail internal-index 125 false true +echo "qualification status regression passed" diff --git a/test/jepsen/qualification-result.sh b/test/jepsen/qualification-result.sh new file mode 100644 index 0000000..e4cdf21 --- /dev/null +++ b/test/jepsen/qualification-result.sh @@ -0,0 +1,19 @@ +#!/usr/bin/env bash + +# Only a normal CLI exit and a fresh, completed checker decision qualify. +# In particular timeout (124/137), timeout invocation (125), and JVM/Docker +# failures are not negative checker evidence, even if an artifact exists. +qualification_result() { + local expectation="$1" corruption="$2" status="$3" result="$4" + local expected_status=1 valid=false + if [[ "${expectation}" == pass ]]; then + expected_status=0 + valid=true + fi + if [[ "${status}" -ne "${expected_status}" ]] || + [[ ! -f "${result}" ]] || + [[ "$(cat "${result}")" != "$(printf 'group-qualification-v1\t%s\t%s\ttrue' "${corruption}" "${valid}")" ]]; then + echo "Jepsen qualification failed: ${expectation}/${corruption}, exit ${status}; missing, mismatched, or unqualified checker result" >&2 + return 1 + fi +} diff --git a/test/jepsen/qualify.sh b/test/jepsen/qualify.sh index 20b96de..4fac84a 100755 --- a/test/jepsen/qualify.sh +++ b/test/jepsen/qualify.sh @@ -4,6 +4,7 @@ set -euo pipefail script_dir="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" repo_dir="$(cd "${script_dir}/../.." && pwd)" artifact_dir="$(mktemp -d "${script_dir}/.cache/qualification.XXXXXX")" +source "${script_dir}/qualification-result.sh" cd "${repo_dir}" @@ -15,10 +16,12 @@ run_jepsen() { local expectation="$1" local corruption="$2" local log="${artifact_dir}/${expectation}-${corruption}.log" + local result="${artifact_dir}/${expectation}-${corruption}.result" local status=0 set +e - timeout --signal=TERM --kill-after=30 180 "${script_dir}/run.sh" test \ + GROUP_JEPSEN_QUALIFICATION_RESULT="${result}" \ + timeout --signal=TERM --kill-after=30 180 "${script_dir}/run.sh" test \ --no-ssh \ --nodes n1,n2,n3 \ --concurrency 2n \ @@ -31,13 +34,8 @@ run_jepsen() { status=$? set -e - if [[ "${expectation}" == "pass" ]] && [[ "${status}" -ne 0 ]]; then - echo "healthy Jepsen baseline failed; see ${log}" >&2 - return 1 - fi - - if [[ "${expectation}" == "fail" ]] && [[ "${status}" -eq 0 ]]; then - echo "Jepsen checker accepted corruption ${corruption}; see ${log}" >&2 + if ! qualification_result "${expectation}" "${corruption}" "${status}" "${result}"; then + echo "see ${log} and ${result}" >&2 return 1 fi diff --git a/test/jepsen/src/group/jepsen/core.clj b/test/jepsen/src/group/jepsen/core.clj index e01ca2d..b6a3f48 100644 --- a/test/jepsen/src/group/jepsen/core.clj +++ b/test/jepsen/src/group/jepsen/core.clj @@ -4,6 +4,7 @@ [group.jepsen.db :as group-db] [group.jepsen.model :as model] [group.jepsen.nemesis :as group-nemesis] + [group.jepsen.qualification :as qualification] [jepsen.cli :as cli] [jepsen.generator :as gen] [jepsen.os :as os] @@ -168,7 +169,7 @@ :nemesis (group-nemesis/nemesis db) :pure-generators true :generator (workload opts) - :checker (model/checker)}))) + :checker (qualification/checker (model/checker))}))) (def cli-options [[nil "--key-count NUMBER" "Number of keys in each cluster and data type" diff --git a/test/jepsen/src/group/jepsen/qualification.clj b/test/jepsen/src/group/jepsen/qualification.clj new file mode 100644 index 0000000..62f875a --- /dev/null +++ b/test/jepsen/src/group/jepsen/qualification.clj @@ -0,0 +1,37 @@ +(ns group.jepsen.qualification + (:require [jepsen.checker :as checker])) + +(defn qualified? [test history result] + (let [mode (keyword (:corruption test)) + injected? (some #(and (= :corrupt (:f %)) + (= :ok (:type %)) + (= mode (get-in % [:value :request :mode]))) + history)] + (boolean + (case mode + :none (true? (:valid? result)) + :unexpected-death (and injected? (seq (:unexpected-owner-deaths result))) + :internal-index (and injected? (seq (:internal-invariant-errors result))) + :cursor-marker (and injected? (seq (:internal-invariant-errors result))) + :registry-projection (and injected? (seq (:internal-invariant-errors result))) + :terminal-unavailable + (let [target (first (:terminal-nodes test))] + (and (some #(and (= :retire-node (:f %)) + (= :info (:type %)) + (= target (get-in % [:value :retired]))) + history) + (contains? (:missing-nodes result) (name target)))) + false)))) + +(defn checker [delegate] + (reify checker/Checker + (check [_ test history opts] + (let [result (checker/check delegate test history opts)] + ;; A dedicated per-run artifact, never human log text. Exceptions and + ;; indeterminate results cannot certify a completed checker decision. + (when-let [path (System/getenv "GROUP_JEPSEN_QUALIFICATION_RESULT")] + (when (boolean? (:valid? result)) + (spit path (str "group-qualification-v1\t" (:corruption test) "\t" + (:valid? result) "\t" + (qualified? test history result) "\n")))) + result)))) diff --git a/test/jepsen/test/group/jepsen/qualification_test.clj b/test/jepsen/test/group/jepsen/qualification_test.clj new file mode 100644 index 0000000..cc5231b --- /dev/null +++ b/test/jepsen/test/group/jepsen/qualification_test.clj @@ -0,0 +1,26 @@ +(ns group.jepsen.qualification-test + (:require [clojure.test :refer :all] + [group.jepsen.qualification :as qualification])) + +(deftest corruption-must-be-injected-and-reach-its-check + (doseq [[mode field] [["unexpected-death" :unexpected-owner-deaths] + ["internal-index" :internal-invariant-errors] + ["cursor-marker" :internal-invariant-errors] + ["registry-projection" :internal-invariant-errors]]] + (let [test {:corruption mode} + history [{:f :corrupt :type :ok + :value {:request {:mode (keyword mode)}}}] + result {:valid? false field {:n1 :evidence}}] + (is (qualification/qualified? test history result)) + (is (not (qualification/qualified? test [] result))) + (is (not (qualification/qualified? test history {:valid? false}))) + (is (not (qualification/qualified? test + [(assoc (first history) :type :fail)] + result)))))) + +(deftest terminal-qualification-requires-retirement-and-missing-target + (let [test {:corruption "terminal-unavailable" :terminal-nodes [:n1 :n2 :n3]} + history [{:f :retire-node :type :info :value {:retired :n1}}]] + (is (qualification/qualified? test history {:missing-nodes #{"n1"}})) + (is (not (qualification/qualified? test [] {:missing-nodes #{"n1"}}))) + (is (not (qualification/qualified? test history {:missing-nodes #{"n2"}}))))) From 598d65eb3bca6b0998ffcb1dc112640a25577d83 Mon Sep 17 00:00:00 2001 From: Jason Stiebs Date: Fri, 11 Sep 2026 12:25:21 -0500 Subject: [PATCH 2/8] Initialize standalone qualification artifacts --- test/jepsen/qualify.sh | 1 + test/jepsen_qualification_cache_test.exs | 55 ++++++++++++++++++++++++ 2 files changed, 56 insertions(+) create mode 100644 test/jepsen_qualification_cache_test.exs diff --git a/test/jepsen/qualify.sh b/test/jepsen/qualify.sh index 20b96de..bbd3c3d 100755 --- a/test/jepsen/qualify.sh +++ b/test/jepsen/qualify.sh @@ -3,6 +3,7 @@ set -euo pipefail script_dir="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" repo_dir="$(cd "${script_dir}/../.." && pwd)" +mkdir -p "${script_dir}/.cache" artifact_dir="$(mktemp -d "${script_dir}/.cache/qualification.XXXXXX")" cd "${repo_dir}" diff --git a/test/jepsen_qualification_cache_test.exs b/test/jepsen_qualification_cache_test.exs new file mode 100644 index 0000000..62de269 --- /dev/null +++ b/test/jepsen_qualification_cache_test.exs @@ -0,0 +1,55 @@ +defmodule Group.JepsenQualificationCacheTest do + use ExUnit.Case, async: true + + @moduletag :local + @moduletag :tmp_dir + + for existing_cache? <- [false, true] do + @tag existing_cache?: existing_cache? + test "standalone qualification prepares artifacts with existing cache: #{existing_cache?}", + %{tmp_dir: directory, existing_cache?: existing_cache?} do + repo = Path.join(directory, "checkout") + script_dir = Path.join(repo, "test/jepsen") + cache = Path.join(script_dir, ".cache") + bin = Path.join(directory, "bin") + File.mkdir_p!(script_dir) + File.mkdir_p!(bin) + File.cp!("test/jepsen/qualify.sh", Path.join(script_dir, "qualify.sh")) + + if existing_cache? do + File.mkdir_p!(cache) + File.write!(Path.join(cache, "retained-artifact"), "previous run") + end + + # Stop at the first external phase, before mutation or Docker work. This + # also verifies that a baseline failure still propagates out of the script. + mix = Path.join(bin, "mix") + + File.write!(mix, """ + #!/bin/sh + printf '%s\\n' "$PWD" "$@" > baseline-invocation + exit 42 + """) + + File.chmod!(mix, 0o755) + + {output, status} = + System.cmd("bash", [Path.join(script_dir, "qualify.sh")], + env: [{"PATH", bin <> ":" <> System.fetch_env!("PATH")}], + stderr_to_stdout: true + ) + + assert status == 42, output + + assert File.read!(Path.join(repo, "baseline-invocation")) == + "#{Path.expand(repo)}\nrun\ntest/mutation/run.exs\n" + + assert [artifact] = Path.wildcard(Path.join(cache, "qualification.*")) + assert File.dir?(artifact) + + if existing_cache? do + assert File.read!(Path.join(cache, "retained-artifact")) == "previous run" + end + end + end +end From b1c1206e810d8b9c50f1b3ac87d06238988a5452 Mon Sep 17 00:00:00 2001 From: Jason Stiebs Date: Fri, 11 Sep 2026 12:27:58 -0500 Subject: [PATCH 3/8] Compare terminal Jepsen metadata with acknowledged owner revisions --- mix.exs | 2 +- test/jepsen/README.md | 14 +++- test/jepsen/metadata_capture.exs | 76 +++++++++++++++++++ test/jepsen/node.exs | 8 +- test/jepsen/project.clj | 2 + test/jepsen/src/group/jepsen/model.clj | 15 ++-- .../group/jepsen/metadata_capture_test.clj | 30 ++++++++ test/jepsen/test/group/jepsen/model_test.clj | 27 ++++++- 8 files changed, 158 insertions(+), 16 deletions(-) create mode 100644 test/jepsen/metadata_capture.exs create mode 100644 test/jepsen/test/group/jepsen/metadata_capture_test.clj diff --git a/mix.exs b/mix.exs index f9b9580..07439f6 100644 --- a/mix.exs +++ b/mix.exs @@ -64,7 +64,7 @@ defmodule Group.MixProject do defp aliases do [ - test: ["test", "cmd test/jepsen/checker.sh"], + test: ["test", "cmd test/jepsen/checker.sh", "cmd test/jepsen/lein.sh test :capture"], "test.soak": [ # Run the PR gate in a child VM. test_helper starts distribution, and # keeping that VM alive for the following `cmd` phases can retain a diff --git a/test/jepsen/README.md b/test/jepsen/README.md index 199c9e7..3770958 100644 --- a/test/jepsen/README.md +++ b/test/jepsen/README.md @@ -32,7 +32,8 @@ The replica lane is selectable without changing the workload or checker: After faults stop, every surviving node reconnects and the harness takes two terminal snapshots. The independent checker requires: -- exact, identical public registry and PG views on every survivor; +- exact, identical public registry and PG metadata on every survivor, including + each revision acknowledged by its owner; - every live owner claim to be visible, and no dead owner token to remain; - deterministic resolution of registry conflicts with no unexpected owner deaths; @@ -135,7 +136,16 @@ test/jepsen/checker.sh ``` At the repository root, `mix test` runs this pure checker after the complete -ExUnit, StreamData, and deterministic-chaos suite. `mix test.soak` runs that +ExUnit, StreamData, and deterministic-chaos suite, followed by executable +capture qualification (`test/jepsen/lein.sh test :capture`, requiring Elixir). +The capture qualification mutates real owners repeatedly, corrupts materialized +metadata while preserving internal index consistency, and passes the real +snapshot EDN to the independent checker. Public values retain metadata maps +rather than only tokens; node/boot/owner/incarnation tokens still identify owner +lifetimes without exporting raw PIDs. Old token-only histories are intentionally +not accepted as exact metadata evidence. + +`mix test.soak` runs that same PR gate, the complete mutation/live-checker qualification, and then `campaign.sh`. Chaos/mixed uses a sender/repair buffer of 32 and requires evidence that one repaired delta run contained at least two records; all other diff --git a/test/jepsen/metadata_capture.exs b/test/jepsen/metadata_capture.exs new file mode 100644 index 0000000..9f1d5db --- /dev/null +++ b/test/jepsen/metadata_capture.exs @@ -0,0 +1,76 @@ +System.put_env("GROUP_JEPSEN_LIBRARY", "1") +Code.require_file("node.exs", __DIR__) + +alias Group.Jepsen.{Driver, EDN, Snapshot} +alias Group.Replica.Data + +{:ok, _} = Group.Jepsen.Transport.Stats.start_link([]) + +{:ok, _} = + Group.start_link( + name: :jepsen_group, + shards: 1, + log: false, + resolve_registry_conflict: {Group.Jepsen.ConflictResolver, :resolve, []} + ) + +{:ok, _} = Group.Jepsen.Driver.Supervisor.start_link(node_id: "n1", boot_id: "capture") + +for revision <- 1..2, operation <- [:register, :join] do + %{status: :ok} = Driver.mutate(operation, "owner", nil, 0, revision) +end + +capture = fn -> Snapshot.capture("n1", "capture", 1, [], []).snapshot end +healthy = capture.() +[owner] = healthy.owners +expected = %{token: owner.token, revision: 2} +^expected = healthy.registry["root"][0] +[^expected] = healthy.pg["root"][0] +[%{revision: 2}] = owner.registrations +[%{revision: 2}] = owner.memberships + +# Corrupt all materialized indexes coherently: the old internal projection +# invariant remains healthy, so only the independent owner/public comparison +# can catch lost, missing or wrong acknowledged metadata. +tables = [ + Data.reg_by_key_table(:jepsen_group, 0), + Data.reg_by_pid_table(:jepsen_group, 0), + Data.reg_claim_by_key_table(:jepsen_group, 0), + Data.reg_claim_by_pid_table(:jepsen_group, 0), + Data.pg_by_key_table(:jepsen_group, 0), + Data.pg_by_pid_table(:jepsen_group, 0) +] + +originals = Map.new(tables, &{&1, :ets.tab2list(&1)}) + +corruptions = + for meta <- [ + %{token: owner.token, revision: 1}, + %{token: owner.token}, + %{token: owner.token, revision: "2"}, + %{token: "wrong-owner", revision: 2}, + %{token: owner.token, revision: 2, extra: true} + ] do + Enum.each(originals, fn {table, rows} -> + changed = + Enum.map(rows, fn row -> + row + |> Tuple.to_list() + |> Enum.map(fn value -> if value == expected, do: meta, else: value end) + |> List.to_tuple() + end) + + :ets.insert(table, changed) + end) + + result = capture.() + true = result.internal.healthy + result + end + +Enum.each(originals, fn {table, rows} -> :ets.insert(table, rows) end) +%{status: :ok} = Driver.kill("owner") +%{status: :ok, owner: reincarnated} = Driver.mutate(:join, "owner", nil, 0, 2) +false = reincarnated.token == owner.token + +File.write!(hd(System.argv()), EDN.encode(%{healthy: healthy, corruptions: corruptions})) diff --git a/test/jepsen/node.exs b/test/jepsen/node.exs index 65bd92a..b2db261 100644 --- a/test/jepsen/node.exs +++ b/test/jepsen/node.exs @@ -1240,7 +1240,7 @@ defmodule Group.Jepsen.Snapshot do value = case Group.lookup(:jepsen_group, registry_key(key), cluster_opts(cluster)) do nil -> nil - {_pid, %{token: token}} -> token + {_pid, %{token: _token} = meta} -> meta {_pid, other} -> "INVALID:#{inspect(other)}" end @@ -1258,7 +1258,7 @@ defmodule Group.Jepsen.Snapshot do :jepsen_group |> Group.members(pg_key(key), cluster_opts(cluster)) |> Enum.map(fn - {_pid, %{token: token}} -> token + {_pid, %{token: _token} = meta} -> meta {_pid, other} -> "INVALID:#{inspect(other)}" end) |> Enum.sort() @@ -1565,4 +1565,6 @@ defmodule Group.Jepsen.Main do end end -Group.Jepsen.Main.run(System.argv()) +unless System.get_env("GROUP_JEPSEN_LIBRARY") == "1" do + Group.Jepsen.Main.run(System.argv()) +end diff --git a/test/jepsen/project.clj b/test/jepsen/project.clj index 05284a5..247c800 100644 --- a/test/jepsen/project.clj +++ b/test/jepsen/project.clj @@ -4,5 +4,7 @@ :license {:name "MIT"} :dependencies [[org.clojure/clojure "1.12.4"] [jepsen "0.3.13"]] + :test-selectors {:default (complement :capture) + :capture :capture} :main group.jepsen.core :jvm-opts ["-Xmx4g" "-Djava.awt.headless=true" "-server"]) diff --git a/test/jepsen/src/group/jepsen/model.clj b/test/jepsen/src/group/jepsen/model.clj index bc772aa..0a95f21 100644 --- a/test/jepsen/src/group/jepsen/model.clj +++ b/test/jepsen/src/group/jepsen/model.clj @@ -40,9 +40,10 @@ (let [owners (->> snapshots vals (mapcat :owners) (map (juxt :token identity)) (into {})) registry-candidates (reduce (fn [by-key [_ owner]] - (reduce (fn [entries {:keys [cluster key]}] + (reduce (fn [entries {:keys [cluster key revision]}] (update entries [(or cluster "root") key] - (fnil conj #{}) (:token owner))) + (fnil conj #{}) {:token (:token owner) + :revision revision})) by-key (owner-entries owner :registrations))) {} @@ -55,9 +56,10 @@ registry-candidates) pg (reduce (fn [view [_ owner]] - (reduce (fn [entries {:keys [cluster key]}] + (reduce (fn [entries {:keys [cluster key revision]}] (update-in entries [(or cluster "root") key] - (fnil conj #{}) (:token owner))) + (fnil conj #{}) {:token (:token owner) + :revision revision})) view (owner-entries owner :memberships))) (empty-view test #{}) @@ -176,11 +178,12 @@ (concat (->> registry vals (mapcat vals) (remove nil?)) (->> pg vals (mapcat vals) (mapcat identity))))) + (map :token) set) expected-tokens (set/union - (->> (:registry expected) vals (mapcat vals) (remove nil?) set) - (->> (:pg expected) vals (mapcat vals) (mapcat identity) set)) + (->> (:registry expected) vals (mapcat vals) (remove nil?) (map :token) set) + (->> (:pg expected) vals (mapcat vals) (mapcat identity) (map :token) set)) orphaned (set/difference actual-tokens live-tokens) missing-live (set/difference expected-tokens actual-tokens) latencies (operation-latencies history) diff --git a/test/jepsen/test/group/jepsen/metadata_capture_test.clj b/test/jepsen/test/group/jepsen/metadata_capture_test.clj new file mode 100644 index 0000000..2b8d430 --- /dev/null +++ b/test/jepsen/test/group/jepsen/metadata_capture_test.clj @@ -0,0 +1,30 @@ +(ns group.jepsen.metadata-capture-test + (:require [clojure.edn :as edn] + [clojure.java.shell :as shell] + [clojure.test :refer :all] + [group.jepsen.model :as model])) + +(def capture-test + {:nodes ["n1"] :key-count 1 :clusters [] :transport "distribution" + :terminal-snapshots-per-node 1 :required-transport-events #{}}) + +(defn analyze-snapshot [snapshot] + (model/analyze capture-test + [{:index 0 :process 0 :type :ok :f :snapshot :value snapshot}])) + +(deftest ^:capture captures-acknowledged-metadata-through-the-real-edn-boundary + (let [output (java.io.File/createTempFile "group-metadata-" ".edn")] + (try + (let [run (shell/sh "env" "ERL_FLAGS=+S 2:2" "MIX_ENV=test" "mix" "run" + "test/jepsen/metadata_capture.exs" (.getPath output) + :dir "../..")] + (is (= 0 (:exit run)) (str (:out run) (:err run))) + (when (zero? (:exit run)) + (let [{:keys [healthy corruptions]} (edn/read-string (slurp output))] + (is (:valid? (analyze-snapshot healthy))) + (doseq [snapshot corruptions] + (let [result (analyze-snapshot snapshot)] + (is (true? (get-in snapshot [:internal :healthy]))) + (is (false? (:valid? result))) + (is (seq (:mismatched-views result)))))))) + (finally (.delete output))))) diff --git a/test/jepsen/test/group/jepsen/model_test.clj b/test/jepsen/test/group/jepsen/model_test.clj index 81e5263..4cd46cb 100644 --- a/test/jepsen/test/group/jepsen/model_test.clj +++ b/test/jepsen/test/group/jepsen/model_test.clj @@ -63,11 +63,29 @@ (defn with-unexpected-death [op token] (assoc-in op [:value :unexpected-deaths] [{:token token, :reason ":boom"}])) +(deftest compares-each-owner-recorded-revision-exactly + (let [owners [(owner "a" [(registration nil 0 2)] [(membership nil 0 3)])] + registry (assoc-in (empty-registry) ["root" 0] {:token "a" :revision 2}) + pg (assoc-in (empty-pg) ["root" 0] [{:token "a" :revision 3}]) + history (mapv #(snapshot-op % (str "n" %) (if (= 1 %) owners []) + registry pg) [1 2 3])] + (is (:valid? (model/analyze test-map history))) + (doseq [field [:registry :pg] + bad [{:token "a" :revision 1} {:token "a"} {:token "b" :revision 2} + {:token "a" :revision "2"} nil]] + (let [changed (assoc-in history [1 :value field "root" 0] + (if (= field :pg) [bad] bad)) + result (model/analyze test-map changed)] + (is (false? (:valid? result))) + (is (contains? (:mismatched-views result) "n2")))))) + (deftest accepts-an-exact-converged-multi-cluster-view (let [owners [(owner "a" [(registration nil 0 1) (registration "red" 1 2)] []) (owner "b" [] [(membership nil 1 2) (membership "red" 0 3)])] - registry (public-view {0 "a", 1 nil} {0 nil, 1 "a"}) - pg (public-view {0 [], 1 ["b"]} {0 ["b"], 1 []}) + registry (public-view {0 {:token "a" :revision 1}, 1 nil} + {0 nil, 1 {:token "a" :revision 2}}) + pg (public-view {0 [], 1 [{:token "b" :revision 2}]} + {0 [{:token "b" :revision 3}], 1 []}) history [(snapshot-op 1 "n1" owners registry pg) (snapshot-op 2 "n2" [] registry pg) (snapshot-op 3 "n3" [] registry pg)]] @@ -96,7 +114,7 @@ (deftest rejects-zombies-missing-live-owners-and-divergence (let [live (owner "live" [(registration nil 0 1)] []) - stale-registry (assoc-in (empty-registry) ["root" 0] "dead") + stale-registry (assoc-in (empty-registry) ["root" 0] {:token "dead" :revision 1}) result (model/analyze test-map [(snapshot-op 1 "n1" [live] stale-registry (empty-pg)) @@ -116,7 +134,8 @@ (snapshot-op 3 "n3" [] registry (empty-pg))] result (model/analyze test-map history)] (is (false? (:valid? result))) - (is (= {["root" 0] #{"a" "b"}} (:live-registry-conflicts result))))) + (is (= {["root" 0] #{{:token "a" :revision 1} {:token "b" :revision 2}}} + (:live-registry-conflicts result))))) (deftest rejects-an-unexpected-owner-death-even-after-cleanup (let [result (model/analyze From 3bb16b370e1f070e0a2766cff31a2fc269de74f9 Mon Sep 17 00:00:00 2001 From: Jason Stiebs Date: Fri, 11 Sep 2026 12:28:50 -0500 Subject: [PATCH 4/8] Isolate mutation targets from the independent checker gate --- mix.exs | 5 +- test/mutation/isolation_test.exs | 139 +++++++++++++++++++++++++++++++ test/mutation/run.exs | 5 +- 3 files changed, 145 insertions(+), 4 deletions(-) create mode 100644 test/mutation/isolation_test.exs diff --git a/mix.exs b/mix.exs index f9b9580..0cc2c3a 100644 --- a/mix.exs +++ b/mix.exs @@ -35,7 +35,7 @@ defmodule Group.MixProject do end def cli do - [preferred_envs: ["test.soak": :test]] + [preferred_envs: ["test.exunit": :test, "test.soak": :test]] end defp deps do @@ -65,6 +65,9 @@ defmodule Group.MixProject do defp aliases do [ test: ["test", "cmd test/jepsen/checker.sh"], + # Mutation targets must measure ExUnit, not the independent JVM gate. + # Invoke the task module directly so the `test` alias is not expanded. + "test.exunit": [&Mix.Tasks.Test.run/1], "test.soak": [ # Run the PR gate in a child VM. test_helper starts distribution, and # keeping that VM alive for the following `cmd` phases can retain a diff --git a/test/mutation/isolation_test.exs b/test/mutation/isolation_test.exs new file mode 100644 index 0000000..50083fe --- /dev/null +++ b/test/mutation/isolation_test.exs @@ -0,0 +1,139 @@ +# Run with: elixir test/mutation/isolation_test.exs +# Exercises the actual campaign against a tiny dependency-free fixture, with +# the real Mix aliases and one real mutation definition. No peers/JVM/Docker. +ExUnit.start() + +defmodule Group.MutationIsolationTest do + use ExUnit.Case, async: false + + @repo Path.expand("../..", __DIR__) + @mutation "drain_oversized_ingress_batch_without_yield" + + setup do + work = Path.join(@repo, "tmp/mutation-isolation-#{System.unique_integer([:positive])}") + File.mkdir_p!(work) + on_exit(fn -> File.rm_rf!(work) end) + + # Keep aliases and preferred environments identical to the real project, + # but remove application code/dependencies unrelated to this harness probe. + ast = @repo |> Path.join("mix.exs") |> File.read!() |> Code.string_to_quoted!() + + {_, functions} = + Macro.prewalk(ast, [], fn + {kind, _, [{name, _, _}, _]} = node, acc + when kind in [:def, :defp] and name in [:aliases, :cli] -> + {node, [Macro.to_string(node) | acc]} + + node, acc -> + {node, acc} + end) + + write(work, "mix.exs", """ + defmodule Isolation.MixProject do + use Mix.Project + def project, do: [app: :isolation, version: "0.0.0", aliases: aliases()] + #{Enum.join(functions, "\n")} + end + """) + + File.mkdir_p!(Path.join(work, "deps")) + File.cp!(Path.join(@repo, "test/mutation/run.exs"), write(work, "test/mutation/run.exs", "")) + write(work, "test/test_helper.exs", "ExUnit.start()\n") + + checker = + write(work, "test/jepsen/checker.sh", """ + #!/bin/sh + echo invoked >> "#{work}/checker-invocations" + exit 73 + """) + + File.chmod!(checker, 0o755) + + write(work, "lib/group/replica.ex", """ + defmodule Isolation.Target do + @incoming_batch_quota 1 + def split(messages) do + {turn, remaining} = Enum.split(messages, @incoming_batch_quota) + {turn, remaining} + end + end + """) + + {:ok, work: work} + end + + test "a passing target survives a broken checker; ordinary mix test still runs it", %{ + work: work + } do + target(work, "assert is_tuple(Isolation.Target.split([1, 2]))") + {output, status} = campaign(work) + assert status == 1, output + assert output =~ "#{@mutation}: SURVIVED" + refute File.exists?(Path.join(work, "checker-invocations")) + + {output, status} = command(work, ["mix", "test", "test/group_test.exs:34"]) + assert status != 0, output + assert File.read!(Path.join(work, "checker-invocations")) == "invoked\n" + end + + test "a regression assertion kills the mutant", %{work: work} do + target(work, "assert Isolation.Target.split([1, 2]) == {[1], [2]}") + {output, status} = campaign(work) + assert status == 0, output + assert output =~ "#{@mutation}: killed" + refute File.exists?(Path.join(work, "checker-invocations")) + end + + test "a failing baseline fails the campaign before mutation", %{work: work} do + target(work, "assert false") + {output, status} = campaign(work) + assert status != 0, output + assert output =~ "baseline failed" + refute output =~ "#{@mutation}: killed" + end + + test "a noncompiling mutant remains invalid, not killed", %{work: work} do + target(work, "assert is_tuple(Isolation.Target.split([1, 2]))") + path = Path.join(work, "test/mutation/run.exs") + source = File.read!(path) + + # Only the fixture's chosen replacement is made syntactically invalid. + old = ~S(" _ = @incoming_batch_quota\n turn = messages\n remaining = []") + assert length(:binary.matches(source, old)) == 1 + File.write!(path, String.replace(source, old, ~S(" this will not compile("))) + + {output, status} = campaign(work) + assert status != 0, output + assert output =~ "#{@mutation}: INVALID (does not compile)" + refute output =~ "#{@mutation}: killed" + end + + defp target(work, assertion) do + write( + work, + "test/group_test.exs", + "defmodule Isolation.TargetTest do\n use ExUnit.Case\n" <> + String.duplicate("\n", 31) <> + " test \"target\" do\n #{assertion}\n end\nend\n" + ) + end + + defp campaign(work), do: command(work, ["elixir", "test/mutation/run.exs", @mutation]) + + defp command(work, args) do + timeout = System.find_executable("timeout") || raise "timeout is required" + + System.cmd(timeout, ["45" | args], + cd: work, + env: [{"ERL_FLAGS", "+S 2:2"}, {"MIX_ENV", nil}, {"ERL_AFLAGS", nil}], + stderr_to_stdout: true + ) + end + + defp write(work, relative, content) do + path = Path.join(work, relative) + File.mkdir_p!(Path.dirname(path)) + File.write!(path, content) + path + end +end diff --git a/test/mutation/run.exs b/test/mutation/run.exs index 96f1f14..12f0a8d 100644 --- a/test/mutation/run.exs +++ b/test/mutation/run.exs @@ -1166,11 +1166,10 @@ defmodule Group.MutationCampaign do end defp run_test(directory, test) do - run_with_timeout(directory, ["mix", "test" | test], + run_with_timeout(directory, ["mix", "test.exunit" | test], env: [ {"GROUP_MODEL_RUNS", "1"}, - {"GROUP_MODEL_COMMANDS", "8"}, - {"GROUP_JEPSEN_SKIP_CHECKER", "1"} + {"GROUP_MODEL_COMMANDS", "8"} ] ) end From 9c2a7fbe57e15689ecdb2c383b9de1d0cb33531a Mon Sep 17 00:00:00 2001 From: Jason Stiebs Date: Fri, 11 Sep 2026 12:33:01 -0500 Subject: [PATCH 5/8] Validate Jepsen registry-conflict deaths against independent claim evidence --- test/jepsen/README.md | 28 +++++ test/jepsen/conflict_probe.exs | 125 +++++++++++++++++++ test/jepsen/node.exs | 97 +++++++++++++- test/jepsen/src/group/jepsen/model.clj | 84 ++++++++++++- test/jepsen/test/group/jepsen/model_test.clj | 77 +++++++++++- 5 files changed, 402 insertions(+), 9 deletions(-) create mode 100644 test/jepsen/conflict_probe.exs diff --git a/test/jepsen/README.md b/test/jepsen/README.md index 199c9e7..10a19d9 100644 --- a/test/jepsen/README.md +++ b/test/jepsen/README.md @@ -55,9 +55,37 @@ cannot masquerade as a current owner. Every history explicitly restarts one node after the deterministic conflict prelude, proving the checker does not mistake restart-sensitive instrumentation for missing protocol coverage. +Registry-conflict exits are obligations, not trusted coverage counters. Before +calling Group, each owner journals its registration attempt independently of +Group's tables; it journals successful or rejected replies, unregisters, and +cluster-intent removal as well. Drivers retain the victim token, key, and +winner metadata from every conflict death. The checker reconstructs the +victim's registrations, requires the reported winner to rank strictly higher +by `{revision, token}`, and requires a matching historical winning claim in +the same cluster and key. Only validated deaths count toward coverage. + +Registration calls interrupted by death or an indeterminate reply remain +possible claims: Group may have installed them before the owner could record +success. A definitive `:taken` reply excludes an attempt. Winning evidence is +retained after unregister, death, and BEAM restart because delayed replicas +can legitimately act on an older claim. Local journal order excludes winners +first attempted after a death; the oracle does not invent a global clock or +infer remote deletion delivery from wall time. This establishes independently +witnessed possible winners, not the exact instant a replica learned a claim. + +The append-only conflict journal retains small operation records for the +bounded campaign, not production ETS rows. Its path defaults to +`/tmp/group-jepsen-conflict-evidence` inside each container and can be overridden +with the driver's `:conflict_evidence_path` option. Corrupt or unreadable +journals fail closed. Checker qualification also loads the real Elixir +Owner/Driver harness, injects valid and forged death reasons, exercises a +register interrupted before its reply, and checks the emitted EDN with the +Clojure lifecycle oracle. + ## Requirements - Docker with Compose v2 +- Elixir/Mix with the repository dependencies installed (checker qualification) - Java 21 or newer - `curl` diff --git a/test/jepsen/conflict_probe.exs b/test/jepsen/conflict_probe.exs new file mode 100644 index 0000000..e15e2a4 --- /dev/null +++ b/test/jepsen/conflict_probe.exs @@ -0,0 +1,125 @@ +# Invoked by the Clojure checker qualification. Load the actual harness without +# starting its TCP server; no copied Owner/Driver logic lives in this probe. +System.put_env("GROUP_JEPSEN_LIBRARY_ONLY", "1") +Code.require_file("node.exs", __DIR__) +Application.ensure_all_started(:group) +Logger.configure(level: :emergency) + +defmodule Group.Jepsen.ConflictProbe do + alias Group.Jepsen.{ConflictEvidence, Driver, EDN} + + def run do + directory = Path.join(__DIR__, ".cache/conflict-probe-#{System.unique_integer([:positive])}") + File.mkdir_p!(directory) + + try do + {winner, rejected, winner_events} = winner(directory) + + cases = [ + {"historical winner subsequently unregistered and died", true, 1, "jepsen/registry/0", + winner, false}, + {"victim killed during register", true, 1, "jepsen/registry/0", winner, true}, + {"forged key and rank", false, 99, "nonexistent", %{token: "invented", revision: -100}, + false}, + {"rightful winner killed", false, 99, "jepsen/registry/0", winner, false}, + {"nonexistent winning claim", false, 1, "jepsen/registry/0", + %{token: "invented", revision: 100}, false}, + {"definitively rejected winning claim", false, 1, "jepsen/registry/0", rejected, false} + ] + + scenarios = + Enum.with_index(cases, fn {label, valid, revision, key, meta, pending}, index -> + events = victim(directory, index, revision, key, meta, pending) + + %{ + label: label, + valid: valid, + snapshots: %{ + "n1" => %{conflict_evidence: events}, + "n2" => %{conflict_evidence: winner_events} + } + } + end) + + IO.puts("CONFLICT-PROBE " <> EDN.encode(scenarios)) + after + File.rm_rf!(directory) + end + end + + defp start(directory, label, node_id) do + {:ok, group} = Group.start_link(name: :jepsen_group, shards: 1, log: false) + path = Path.join(directory, label) + {:ok, evidence} = ConflictEvidence.start_link(conflict_evidence_path: path) + {:ok, driver} = Driver.start_link(index: 0, node_id: node_id, boot_id: label) + {group, evidence, driver, path} + end + + defp mutate(driver, operation, revision) do + GenServer.call(driver, {:mutate, operation, "owner", nil, 0, revision}) + end + + defp winner(directory) do + {group, evidence, driver, _path} = start(directory, "winner", "n2") + %{status: :ok, owner: %{token: token}} = mutate(driver, :register, 10) + + %{status: :fail, owner: %{token: rejected}} = + GenServer.call(driver, {:mutate, :register, "rejected", nil, 0, 100}) + + %{status: :ok} = mutate(driver, :unregister, 0) + %{status: :ok} = GenServer.call(driver, {:kill, "owner"}) + %{status: :ok} = GenServer.call(driver, {:kill, "rejected"}) + events = ConflictEvidence.snapshot() + GenServer.stop(driver) + GenServer.stop(evidence) + Supervisor.stop(group) + {%{token: token, revision: 10}, %{token: rejected, revision: 100}, events} + end + + defp victim(directory, index, revision, key, winner, pending) do + {group, evidence, driver, path} = start(directory, "victim-#{index}", "n1") + %{status: :ok} = mutate(driver, :join, 0) + {pid, _token, _monitor, _cached} = :sys.get_state(driver).owners["owner"] + shard = Group.Replica.shard_for(:jepsen_group, nil, "jepsen/registry/0") + + task = + if pending do + :ok = :sys.suspend(shard) + task = Task.async(fn -> mutate(driver, :register, revision) end) + wait(fn -> Enum.any?(ConflictEvidence.snapshot(), &(&1.kind == :register)) end) + task + else + %{status: :ok} = mutate(driver, :register, revision) + nil + end + + Process.exit(pid, {:group_registry_conflict, key, winner}) + if task, do: Task.await(task) + wait(fn -> Enum.any?(ConflictEvidence.snapshot(), &(&1.kind == :death)) end) + if pending, do: :sys.resume(shard) + events = ConflictEvidence.snapshot() + [] = GenServer.call(driver, :unexpected_deaths) + {:ok, []} = GenServer.call(driver, :owner_snapshots) + GenServer.stop(driver) + GenServer.stop(evidence) + + # The exact registration/death obligations must survive a recorder restart. + {:ok, evidence} = ConflictEvidence.start_link(conflict_evidence_path: path) + ^events = ConflictEvidence.snapshot() + GenServer.stop(evidence) + Supervisor.stop(group) + events + end + + defp wait(fun, remaining \\ 200) + defp wait(_fun, 0), do: raise("probe timed out") + + defp wait(fun, remaining) do + unless fun.() do + Process.sleep(5) + wait(fun, remaining - 1) + end + end +end + +Group.Jepsen.ConflictProbe.run() diff --git a/test/jepsen/node.exs b/test/jepsen/node.exs index 65bd92a..152cc96 100644 --- a/test/jepsen/node.exs +++ b/test/jepsen/node.exs @@ -424,6 +424,51 @@ defmodule Group.Jepsen.ConflictResolver do defp rank(_meta), do: {-1, ""} end +defmodule Group.Jepsen.ConflictEvidence do + @moduledoc false + use GenServer + + # Independent of Group's ETS, and retained across container/BEAM restarts. + # Persist the invocation before calling Group: a conflict can kill an owner + # before its register call returns, even though its claim was installed. + def start_link(opts), do: GenServer.start_link(__MODULE__, opts, name: __MODULE__) + def record(event), do: GenServer.call(__MODULE__, {:record, event}) + def snapshot, do: GenServer.call(__MODULE__, :snapshot) + + @impl true + def init(opts) do + path = Keyword.get(opts, :conflict_evidence_path, "/tmp/group-jepsen-conflict-evidence") + + events = + case File.read(path) do + {:ok, contents} -> + contents + |> String.split("\n", trim: true) + |> Enum.map(&(&1 |> Base.decode64!() |> :erlang.binary_to_term())) + + {:error, :enoent} -> + [] + + {:error, reason} -> + raise "cannot read conflict evidence: #{inspect(reason)}" + end + + {:ok, %{path: path, events: Enum.reverse(events), sequence: length(events)}} + end + + @impl true + def handle_call({:record, event}, _from, state) do + event = Map.put(event, :sequence, state.sequence + 1) + encoded = event |> :erlang.term_to_binary() |> Base.encode64() + :ok = File.write(state.path, encoded <> "\n", [:append, :sync]) + {:reply, event.sequence, %{state | events: [event | state.events], sequence: event.sequence}} + end + + def handle_call(:snapshot, _from, state) do + {:reply, Enum.reverse(state.events), state} + end +end + defmodule Group.Jepsen.Owner do @moduledoc false use GenServer @@ -437,15 +482,26 @@ defmodule Group.Jepsen.Owner do def handle_call({:mutate, :register, cluster, key, revision}, _from, state) do meta = %{token: state.token, revision: revision} + attempt = + Group.Jepsen.ConflictEvidence.record(%{ + kind: :register, + token: state.token, + cluster: cluster, + key: key, + revision: revision + }) + case safe_group_call(fn -> Group.register(:jepsen_group, registry_key(key), meta, cluster_opts(cluster)) end) do :ok -> + registration_result(attempt, :ok) entry = %{cluster: cluster, key: key, revision: revision} state = put_in(state.registrations[{cluster, key}], entry) {:reply, {:ok, snapshot(state)}, state} {:error, reason} -> + registration_result(attempt, if(reason == :taken, do: :fail, else: :unknown)) {:reply, {:error, reason, snapshot(state)}, state} end end @@ -458,6 +514,13 @@ defmodule Group.Jepsen.Owner do Group.unregister(:jepsen_group, registry_key(key), cluster_opts(cluster)) end) do :ok -> + Group.Jepsen.ConflictEvidence.record(%{ + kind: :unregister, + token: state.token, + cluster: cluster, + key: key + }) + state = %{state | registrations: Map.delete(state.registrations, owner_key)} {:reply, {:ok, snapshot(state)}, state} @@ -505,6 +568,12 @@ defmodule Group.Jepsen.Owner do end def handle_call({:drop_cluster, cluster}, _from, state) do + Group.Jepsen.ConflictEvidence.record(%{ + kind: :drop_cluster, + token: state.token, + cluster: cluster + }) + registrations = drop_cluster(state.registrations, cluster) memberships = drop_cluster(state.memberships, cluster) state = %{state | registrations: registrations, memberships: memberships} @@ -513,6 +582,10 @@ defmodule Group.Jepsen.Owner do def handle_call(:snapshot, _from, state), do: {:reply, snapshot(state), state} + defp registration_result(attempt, status) do + Group.Jepsen.ConflictEvidence.record(%{kind: :result, attempt: attempt, status: status}) + end + defp drop_cluster(entries, cluster) do entries |> Enum.reject(fn {{entry_cluster, _key}, _entry} -> entry_cluster == cluster end) @@ -547,8 +620,6 @@ defmodule Group.Jepsen.Driver do @moduledoc false use GenServer - alias Group.Jepsen.Transport.Stats - @driver_count 8 @unexpected_death_log "/tmp/group-jepsen-unexpected-deaths" @@ -731,7 +802,16 @@ defmodule Group.Jepsen.Driver do unexpected_deaths = if match?({:group_registry_conflict, _key, _winner_meta}, reason) do - Stats.increment_persistent(:registry_conflict_death) + {:group_registry_conflict, key, winner_meta} = reason + + Group.Jepsen.ConflictEvidence.record(%{ + kind: :death, + token: token, + key: if(is_binary(key), do: key, else: %{invalid: inspect(key)}), + winner: + if(is_map(winner_meta), do: winner_meta, else: %{invalid: inspect(winner_meta)}) + }) + state.unexpected_deaths else death = %{token: token, reason: inspect(reason)} @@ -821,7 +901,11 @@ defmodule Group.Jepsen.Driver.Supervisor do @impl true def init(opts), - do: Supervisor.init(Group.Jepsen.Driver.child_specs(opts), strategy: :one_for_one) + do: + Supervisor.init( + [{Group.Jepsen.ConflictEvidence, opts} | Group.Jepsen.Driver.child_specs(opts)], + strategy: :one_for_one + ) end defmodule Group.Jepsen.Cluster do @@ -1277,6 +1361,7 @@ defmodule Group.Jepsen.Snapshot do peers: Group.nodes(:jepsen_group) |> Enum.map(&Atom.to_string/1) |> Enum.sort(), owners: owners, unexpected_deaths: Group.Jepsen.Driver.unexpected_deaths(), + conflict_evidence: Group.Jepsen.ConflictEvidence.snapshot(), transport_events: Group.Jepsen.Transport.Stats.snapshot(), transport_profile: Group.Jepsen.Transport.Control.profile(), internal: Group.Jepsen.Invariant.snapshot(retired_nodes), @@ -1565,4 +1650,6 @@ defmodule Group.Jepsen.Main do end end -Group.Jepsen.Main.run(System.argv()) +unless System.get_env("GROUP_JEPSEN_LIBRARY_ONLY") == "1" do + Group.Jepsen.Main.run(System.argv()) +end diff --git a/test/jepsen/src/group/jepsen/model.clj b/test/jepsen/src/group/jepsen/model.clj index bc772aa..3082705 100644 --- a/test/jepsen/src/group/jepsen/model.clj +++ b/test/jepsen/src/group/jepsen/model.clj @@ -88,6 +88,7 @@ {:owners (set (:owners snapshot)) :peers (set (:peers snapshot)) :unexpected-deaths (set (:unexpected-deaths snapshot)) + :conflict-evidence (:conflict-evidence snapshot) :view (normalize-view test snapshot) :internal (stable-internal snapshot)}) @@ -96,6 +97,79 @@ (remove history/invoke?) (keep #(get-in % [:value :response :latency-us])))) +(defn replay-conflict-evidence [node events] + (reduce + (fn [state {:keys [kind sequence token cluster key attempt status] :as event}] + (let [slot [cluster key]] + (case kind + :register (-> state + (assoc-in [:claims sequence] (assoc event :node node)) + (assoc-in [:pending token sequence] slot)) + :result (let [claim (get-in state [:claims attempt]) + token (:token claim) + slot [(:cluster claim) (:key claim)]] + (cond-> (assoc-in state [:claims attempt :status] status) + (not= :unknown status) (update-in [:pending token] dissoc attempt) + (= :ok status) (assoc-in [:active token slot] attempt))) + :unregister + (-> state + (update-in [:active token] dissoc slot) + (update-in [:pending token] + #(into {} (remove (fn [[_ s]] (= slot s))) %))) + :drop-cluster + (-> state + (update-in [:active token] + #(into {} (remove (fn [[[c _] _]] (= cluster c))) %)) + (update-in [:pending token] + #(into {} (remove (fn [[_ [c _]]] (= cluster c))) %))) + :death + (let [attempts (concat (vals (get-in state [:active token])) + (keys (get-in state [:pending token])))] + (-> state + (update :deaths conj + (assoc event :node node :victim-attempts (set attempts))) + (update :active dissoc token) + (update :pending dissoc token))) + state))) + {:claims {} :active {} :pending {} :deaths []} + events)) + +(defn conflict-analysis [snapshots] + (let [journals (into {} (map (fn [[node snapshot]] + [node (replay-conflict-evidence + node (:conflict-evidence snapshot))])) + snapshots) + claims (mapcat (comp vals :claims val) journals) + possible? #(not= :fail (:status %)) + winners (group-by (juxt :token :revision :key) (filter possible? claims)) + deaths (mapcat (comp :deaths val) journals) + justified? + (fn [{:keys [node sequence token key winner victim-attempts]}] + (boolean + (some + (fn [id] + (let [victim (get-in journals [node :claims id]) + rank (juxt :revision :token)] + (and (possible? victim) + (= token (:token victim)) + (= key (str "jepsen/registry/" (:key victim))) + (integer? (:revision winner)) + (string? (:token winner)) + (not= token (:token winner)) + (pos? (compare (rank winner) (rank victim))) + (some #(and (= (:cluster victim) (:cluster %)) + ;; Local order is known. Across nodes there + ;; is no synchronized oracle clock; delayed + ;; remote deletions may still lose a conflict. + (or (not= node (:node %)) + (< (:sequence %) sequence))) + (get winners [(:token winner) (:revision winner) + (:key victim)]))))) + victim-attempts))) + invalid (vec (remove justified? deaths))] + {:invalid invalid + :validated-count (- (count deaths) (count invalid))})) + (defn analyze [test history] (let [observations (snapshots-by-node history) snapshots (latest-snapshots history) @@ -137,10 +211,12 @@ (when (not= expected-peers actual-peers) [node {:expected expected-peers, :actual actual-peers}])))) relevant-snapshots) + conflicts (conflict-analysis relevant-snapshots) transport-events - (reduce #(merge-with + %1 %2) - {} - (map #(or (:transport-events %) {}) (vals relevant-snapshots))) + (assoc (reduce #(merge-with + %1 %2) + {} + (map #(or (:transport-events %) {}) (vals relevant-snapshots))) + :registry-conflict-death (:validated-count conflicts)) required-transport-events (get test :required-transport-events default-required-transport-events) @@ -198,6 +274,7 @@ (empty? (:conflicts expected)) (empty? mismatches) (empty? unexpected-deaths) + (empty? (:invalid conflicts)) (empty? orphaned) (empty? missing-live) (not latency-violation?))] @@ -219,6 +296,7 @@ :live-registry-conflicts (:conflicts expected) :mismatched-views mismatches :unexpected-owner-deaths unexpected-deaths + :invalid-conflict-deaths (:invalid conflicts) :orphaned-owner-tokens orphaned :missing-live-owner-tokens missing-live :expected expected-view})) diff --git a/test/jepsen/test/group/jepsen/model_test.clj b/test/jepsen/test/group/jepsen/model_test.clj index 81e5263..440fe2b 100644 --- a/test/jepsen/test/group/jepsen/model_test.clj +++ b/test/jepsen/test/group/jepsen/model_test.clj @@ -1,5 +1,8 @@ (ns group.jepsen.model-test (:require [clojure.test :refer :all] + [clojure.edn :as edn] + [clojure.java.shell :as shell] + [clojure.string :as str] [group.jepsen.model :as model])) (def test-map @@ -63,6 +66,78 @@ (defn with-unexpected-death [op token] (assoc-in op [:value :unexpected-deaths] [{:token token, :reason ":boom"}])) +(def conflict-journal + [{:kind :register :sequence 1 :token "a" :cluster nil :key 0 :revision 1} + {:kind :result :sequence 2 :attempt 1 :status :ok} + {:kind :register :sequence 3 :token "b" :cluster nil :key 0 :revision 2} + {:kind :death :sequence 4 :token "a" :key "jepsen/registry/0" + :winner {:token "b" :revision 2}} + {:kind :result :sequence 5 :attempt 3 :status :ok} + {:kind :unregister :sequence 6 :token "b" :cluster nil :key 0}]) + +(defn with-conflict-evidence [op] + (assoc-in op [:value :conflict-evidence] conflict-journal)) + +(deftest validates-historical-conflicts-independently + (let [check #(model/conflict-analysis {"n1" {:conflict-evidence %}}) + valid #(is (empty? (:invalid (check %)))) + invalid #(is (= 1 (count (:invalid (check %)))))] + (valid conflict-journal) + (valid (vec (remove #(= 2 (:sequence %)) conflict-journal))) + (invalid (assoc-in conflict-journal [3 :key] "nonexistent")) + (invalid (assoc-in conflict-journal [3 :winner :revision] -100)) + (invalid (assoc-in conflict-journal [3 :winner :token] "invented")) + (invalid (assoc-in conflict-journal [3 :winner :token] "a")) + (invalid (assoc-in conflict-journal [2 :cluster] "red")) + (invalid (assoc-in conflict-journal [4 :status] :fail)) + (invalid (assoc-in conflict-journal [0 :revision] 99)) + (invalid (vec (concat (subvec conflict-journal 0 3) + [{:kind :unregister :sequence 4 :token "a" :cluster nil :key 0}] + (subvec conflict-journal 3)))) + (invalid (vec (concat (subvec conflict-journal 0 3) + [{:kind :drop-cluster :sequence 4 :token "a" :cluster nil}] + (subvec conflict-journal 3)))) + (invalid (conj conflict-journal + {:kind :register :sequence 7 :token "b" :cluster nil :key 0 :revision 2} + {:kind :death :sequence 8 :token "b" :key "jepsen/registry/0" + :winner {:token "a" :revision 1}})) + (invalid (vec (remove #(= :register (:kind %)) conflict-journal))) + (invalid (vec (concat (subvec conflict-journal 0 2) + [(last conflict-journal)] + (subvec conflict-journal 3 4) + [(assoc (get conflict-journal 2) :sequence 5)]))) + ;; Unique incarnation tokens break equal-revision ties without pid order. + (valid (assoc-in conflict-journal [0 :revision] 2)))) + +(deftest conflict-statistics-alone-do-not-discharge-deaths + (let [history [(assoc-in (snapshot-op 1 "n1" [] (empty-registry) (empty-pg)) + [:value :transport-events] {:registry-conflict-death 99}) + (snapshot-op 2 "n2" [] (empty-registry) (empty-pg)) + (snapshot-op 3 "n3" [] (empty-registry) (empty-pg))] + result (model/analyze (assoc test-map :required-transport-events + #{:registry-conflict-death}) history)] + (is (false? (:valid? result))) + (is (= #{:registry-conflict-death} (:missing-transport-events result))))) + +(deftest checks-real-owner-and-driver-evidence + (let [{:keys [exit out err]} + (shell/sh "mix" "run" "--no-start" "test/jepsen/conflict_probe.exs" + :dir "../.." + :env (assoc (into {} (System/getenv)) "ERL_FLAGS" "+S 2:2")) + _ (is (zero? exit) (str out err)) + line (first (filter #(str/starts-with? % "CONFLICT-PROBE ") + (str/split-lines out))) + scenarios (when line (edn/read-string (subs line 15)))] + (is (some? scenarios) (str out err)) + (doseq [{:keys [label valid snapshots]} scenarios] + (let [history (mapv (fn [index node] + (assoc-in (snapshot-op index node [] (empty-registry) (empty-pg)) + [:value :conflict-evidence] + (get-in snapshots [node :conflict-evidence]))) + (range 3) ["n1" "n2" "n3"]) + result (model/analyze test-map history)] + (is (= valid (:valid? result)) (str label ": " result)))))) + (deftest accepts-an-exact-converged-multi-cluster-view (let [owners [(owner "a" [(registration nil 0 1) (registration "red" 1 2)] []) (owner "b" [] [(membership nil 1 2) (membership "red" 0 3)])] @@ -173,7 +248,7 @@ :snapshot-chunk 2 :multi-chunk-snapshot 1 :registry-conflict-death 1} - history [(assoc-in (snapshot-op 1 "n1" [] (empty-registry) (empty-pg)) + history [(assoc-in (with-conflict-evidence (snapshot-op 1 "n1" [] (empty-registry) (empty-pg))) [:value :transport-events] events) (snapshot-op 2 "n2" [] (empty-registry) (empty-pg)) From b7768f7e138c16c1939d24afab0cde7e3f20c5ba Mon Sep 17 00:00:00 2001 From: Jason Stiebs Date: Fri, 11 Sep 2026 12:42:08 -0500 Subject: [PATCH 6/8] Require target-specific injection and invariant evidence --- test/jepsen/invariant_qualification_test.exs | 90 +++++++++++++++++++ test/jepsen/node.exs | 57 +++++++++--- .../jepsen/src/group/jepsen/qualification.clj | 28 +++++- .../test/group/jepsen/qualification_test.clj | 44 ++++++++- 4 files changed, 199 insertions(+), 20 deletions(-) create mode 100644 test/jepsen/invariant_qualification_test.exs diff --git a/test/jepsen/invariant_qualification_test.exs b/test/jepsen/invariant_qualification_test.exs new file mode 100644 index 0000000..64eb1f8 --- /dev/null +++ b/test/jepsen/invariant_qualification_test.exs @@ -0,0 +1,90 @@ +# MIX_ENV=test mix run --no-start test/jepsen/invariant_qualification_test.exs +# Compile only the actual oracle modules, not the node entrypoint. Redirect +# the injection marker into the workspace so this probe needs no Docker or +# machine-global files. +ExUnit.start() + +defmodule Group.Jepsen.InvariantQualificationTest do + use ExUnit.Case, async: false + + alias Group.Jepsen.Invariant + alias Group.Replica.Data + + @compile {:no_warn_undefined, [Group.Jepsen.Invariant, Group.Jepsen.EDN]} + + setup_all do + work = + Path.expand( + "../../tmp/invariant-qualification-#{System.unique_integer([:positive])}", + __DIR__ + ) + + File.mkdir_p!(work) + marker = Path.join(work, "cursor-marker") + on_exit(fn -> File.rm_rf!(work) end) + {:ok, _} = Application.ensure_all_started(:group) + + {:__block__, _, expressions} = + __DIR__ + |> Path.join("node.exs") + |> File.read!() + |> Code.string_to_quoted!() + + modules = + Enum.filter(expressions, fn + {:defmodule, _, [{:__aliases__, _, [:Group, :Jepsen, name]}, _]} -> + name in [:InvariantViolation, :Invariant, :ConflictResolver, :EDN] + + _ -> + false + end) + + ast = + Macro.prewalk({:__block__, [], modules}, fn + "/tmp/group-jepsen-cursor-marker-corruption" -> marker + node -> node + end) + + Code.compile_quoted(ast) + {:ok, marker: marker} + end + + setup %{marker: marker} do + File.rm(marker) + start_supervised!({Group, name: :jepsen_group, shards: 1, log: false}) + :ok + end + + test "arming with no remote cursor reports no injection or invariant evidence", %{ + marker: marker + } do + File.write!(marker, "enabled\n") + snapshot = Invariant.snapshot([]) + refute snapshot.healthy + assert snapshot.snapshot_staging_count == -1 + assert snapshot.injected_corruptions == [] + assert snapshot.failed_invariants == [] + end + + test "a real marker insertion reports its exact invariant", %{marker: marker} do + File.write!(marker, "enabled\n") + :ets.insert(Data.replica_cursor_table(:jepsen_group, 0), {:probe_stream, 1}) + snapshot = Invariant.snapshot([]) + refute snapshot.healthy + assert snapshot.injected_corruptions == [:"cursor-marker"] + assert snapshot.failed_invariants == [:cursor_snapshot_marker] + assert Group.Jepsen.EDN.encode(snapshot) =~ ":cursor-snapshot-marker" + end + + test "an unrelated index failure cannot masquerade as a cursor marker" do + :ets.insert( + Data.reg_by_pid_table(:jepsen_group, 0), + {{self(), nil, "qualification-probe"}, %{}, 0, node()} + ) + + snapshot = Invariant.snapshot([]) + refute snapshot.healthy + assert snapshot.injected_corruptions == [] + assert snapshot.failed_invariants == [:registry_dual_indexes] + end +end diff --git a/test/jepsen/node.exs b/test/jepsen/node.exs index 65bd92a..4d8253d 100644 --- a/test/jepsen/node.exs +++ b/test/jepsen/node.exs @@ -860,6 +860,10 @@ defmodule Group.Jepsen.Cluster do defp result({:error, reason}), do: %{status: :fail, error: inspect(reason)} end +defmodule Group.Jepsen.InvariantViolation do + defexception [:message, :invariant] +end + defmodule Group.Jepsen.Invariant do @moduledoc false @@ -868,9 +872,9 @@ defmodule Group.Jepsen.Invariant do def snapshot(retired_nodes) do config = Group.get_config(:jepsen_group) shards = 0..(config.num_shards - 1) - maybe_inject_cursor_marker_corruption(shards) + injected_corruptions = maybe_inject_cursor_marker_corruption(shards) - errors = + failures = check("dual indexes", &assert_dual_indexes/0) ++ check("registry claims", &assert_registry_claims/0) ++ check("oplog", &assert_oplogs/0) ++ @@ -889,8 +893,10 @@ defmodule Group.Jepsen.Invariant do end) %{ - healthy: errors == [] and staging_count == 0, - errors: errors, + healthy: failures == [] and staging_count == 0, + errors: Enum.map(failures, & &1.message), + failed_invariants: Enum.flat_map(failures, &List.wrap(&1.invariant)), + injected_corruptions: injected_corruptions, snapshot_staging_count: staging_count, oplog_entries: oplog_entries, oplog_max_entries_per_shard: config.replicated_oplog_max_entries, @@ -905,6 +911,8 @@ defmodule Group.Jepsen.Invariant do %{ healthy: false, errors: ["invariant snapshot failed: #{Exception.message(exception)}"], + failed_invariants: [], + injected_corruptions: [], snapshot_staging_count: -1 } end @@ -928,9 +936,13 @@ defmodule Group.Jepsen.Invariant do {stream, {:snapshot_installing, 1}} ) + [:"cursor-marker"] + nil -> raise "no remote replica cursor available for corruption" end + else + [] end end @@ -938,9 +950,13 @@ defmodule Group.Jepsen.Invariant do fun.() [] rescue - exception -> ["#{label}: #{Exception.message(exception)}"] + exception in Group.Jepsen.InvariantViolation -> + [%{invariant: exception.invariant, message: "#{label}: #{Exception.message(exception)}"}] + + exception -> + [%{invariant: nil, message: "#{label}: #{Exception.message(exception)}"}] catch - kind, reason -> ["#{label}: #{inspect({kind, reason})}"] + kind, reason -> [%{invariant: nil, message: "#{label}: #{inspect({kind, reason})}"}] end defp assert_dual_indexes do @@ -973,7 +989,13 @@ defmodule Group.Jepsen.Invariant do {cluster, key, pid, meta, time, origin} end) - assert_equal!(reg_key, reg_pid, "registry dual indexes shard #{shard}") + assert_equal!( + reg_key, + reg_pid, + "registry dual indexes shard #{shard}", + :registry_dual_indexes + ) + assert_equal!(pg_key, pg_pid, "PG dual indexes shard #{shard}") expected_counts = @@ -1074,7 +1096,7 @@ defmodule Group.Jepsen.Invariant do {cluster, key, pid, meta, time, origin} end) - assert_equal!(expected, visible, "registry projection shard #{shard}") + assert_equal!(expected, visible, "registry projection shard #{shard}", :registry_projection) end) end @@ -1140,6 +1162,12 @@ defmodule Group.Jepsen.Invariant do Data.replica_cursor_table(:jepsen_group, shard) |> :ets.tab2list() |> Enum.each(fn {stream, seq} -> + if match?({:snapshot_installing, _}, seq) do + raise Group.Jepsen.InvariantViolation, + invariant: :cursor_snapshot_marker, + message: "cursor contains uncommitted snapshot marker #{inspect({stream, seq})}" + end + origin = WireProtocol.stream_origin(stream) cluster = WireProtocol.stream_cluster(stream) @@ -1202,10 +1230,13 @@ defmodule Group.Jepsen.Invariant do Enum.each(0..(num_shards - 1), fun) end - defp assert_equal!(left, right, label) do + defp assert_equal!(left, right, label, invariant \\ nil) do if left != right do - raise "#{label}: left-only=#{inspect(MapSet.difference(left, right))} " <> - "right-only=#{inspect(MapSet.difference(right, left))}" + raise Group.Jepsen.InvariantViolation, + invariant: invariant, + message: + "#{label}: left-only=#{inspect(MapSet.difference(left, right))} " <> + "right-only=#{inspect(MapSet.difference(right, left))}" end end @@ -1439,7 +1470,7 @@ defmodule Group.Jepsen.Wire do defp corrupt("internal-index") do table = Group.Replica.Data.reg_by_pid_table(:jepsen_group, 0) :ets.insert(table, {{self(), nil, "jepsen/registry/corrupt"}, %{}, 0, node()}) - %{status: :ok} + %{status: :ok, injected: :"internal-index"} end defp corrupt("cursor-marker") do @@ -1483,7 +1514,7 @@ defmodule Group.Jepsen.Wire do end) if corrupted do - %{status: :ok} + %{status: :ok, injected: :"registry-projection"} else %{status: :fail, error: "no visible registry claim available for corruption"} end diff --git a/test/jepsen/src/group/jepsen/qualification.clj b/test/jepsen/src/group/jepsen/qualification.clj index 62f875a..fbe2cdb 100644 --- a/test/jepsen/src/group/jepsen/qualification.clj +++ b/test/jepsen/src/group/jepsen/qualification.clj @@ -1,6 +1,28 @@ (ns group.jepsen.qualification (:require [jepsen.checker :as checker])) +(def internal-invariants + {:internal-index :registry-dual-indexes + :cursor-marker :cursor-snapshot-marker + :registry-projection :registry-projection}) + +(defn internal-qualified? [mode history result] + ;; Match completion and the precise assertion on the same injected node. + ;; Arming the cursor marker file is not injection: only a successful snapshot + ;; insertion can attest that a remote cursor existed and was corrupted. + (some (fn [op] + (let [internal (get-in result [:internal-invariant-errors + (get-in op [:value :node])]) + completed? (if (= :cursor-marker mode) + (some #{mode} (:injected-corruptions internal)) + (= mode (get-in op [:value :response :injected])))] + (and (= :corrupt (:f op)) + (= :ok (:type op)) + (= mode (get-in op [:value :request :mode])) + completed? + (some #{(internal-invariants mode)} (:failed-invariants internal))))) + history)) + (defn qualified? [test history result] (let [mode (keyword (:corruption test)) injected? (some #(and (= :corrupt (:f %)) @@ -11,9 +33,9 @@ (case mode :none (true? (:valid? result)) :unexpected-death (and injected? (seq (:unexpected-owner-deaths result))) - :internal-index (and injected? (seq (:internal-invariant-errors result))) - :cursor-marker (and injected? (seq (:internal-invariant-errors result))) - :registry-projection (and injected? (seq (:internal-invariant-errors result))) + :internal-index (internal-qualified? mode history result) + :cursor-marker (internal-qualified? mode history result) + :registry-projection (internal-qualified? mode history result) :terminal-unavailable (let [target (first (:terminal-nodes test))] (and (some #(and (= :retire-node (:f %)) diff --git a/test/jepsen/test/group/jepsen/qualification_test.clj b/test/jepsen/test/group/jepsen/qualification_test.clj index cc5231b..f1dfa80 100644 --- a/test/jepsen/test/group/jepsen/qualification_test.clj +++ b/test/jepsen/test/group/jepsen/qualification_test.clj @@ -3,10 +3,7 @@ [group.jepsen.qualification :as qualification])) (deftest corruption-must-be-injected-and-reach-its-check - (doseq [[mode field] [["unexpected-death" :unexpected-owner-deaths] - ["internal-index" :internal-invariant-errors] - ["cursor-marker" :internal-invariant-errors] - ["registry-projection" :internal-invariant-errors]]] + (doseq [[mode field] [["unexpected-death" :unexpected-owner-deaths]]] (let [test {:corruption mode} history [{:f :corrupt :type :ok :value {:request {:mode (keyword mode)}}}] @@ -18,6 +15,45 @@ [(assoc (first history) :type :fail)] result)))))) +(deftest internal-corruption-requires-specific-completion-and-invariant + (doseq [[mode invariant] qualification/internal-invariants] + (let [test {:corruption (name mode)} + op {:f :corrupt :type :ok + :value {:node "n1" :request {:mode mode} + :response {:injected mode}}} + internal {:failed-invariants [invariant] :injected-corruptions [mode]} + result {:valid? false :internal-invariant-errors {"n1" internal}} + qualifies? #(qualification/qualified? test [%1] %2)] + (is (qualifies? op result)) + (is (not (qualifies? (assoc op :type :fail) result))) + (is (not (qualifies? (assoc-in op [:value :node] "n2") result))) + (is (not (qualifies? op (assoc-in result + [:internal-invariant-errors "n1" :failed-invariants] + [:unrelated-invariant])))) + (is (not (qualifies? (update-in op [:value] dissoc :response) + (assoc-in result + [:internal-invariant-errors "n1" :injected-corruptions] + []))))))) + +(deftest arming-cursor-injection-with-no-remote-cursor-is-not-qualification + (let [test {:corruption "cursor-marker"} + history [{:f :corrupt :type :ok + :value {:node "n1" :request {:mode :cursor-marker} + :response {:status :ok}}}] + result {:valid? false + :internal-invariant-errors + {"n1" {:healthy false + :errors ["invariant snapshot failed: no remote replica cursor available for corruption"] + :failed-invariants [] + :injected-corruptions [] + :snapshot-staging-count -1}}}] + (is (not (qualification/qualified? test history result))) + ;; Even a matching assertion elsewhere cannot replace injection completion. + (is (not (qualification/qualified? + test history + (assoc-in result [:internal-invariant-errors "n1" :failed-invariants] + [:cursor-snapshot-marker])))))) + (deftest terminal-qualification-requires-retirement-and-missing-target (let [test {:corruption "terminal-unavailable" :terminal-nodes [:n1 :n2 :n3]} history [{:f :retire-node :type :info :value {:retired :n1}}]] From 7b88453b79bd13726881cfd10c3d2bd3ba8e1326 Mon Sep 17 00:00:00 2001 From: Jason Stiebs Date: Fri, 11 Sep 2026 16:13:53 -0500 Subject: [PATCH 7/8] Archive retired conflict evidence and reset it between histories --- test/jepsen/README.md | 13 ++++- test/jepsen/conflict_probe.exs | 31 +++++++++-- test/jepsen/decode_conflict_evidence.exs | 6 +++ test/jepsen/node.exs | 32 ++++++++++-- test/jepsen/src/group/jepsen/db.clj | 4 ++ test/jepsen/src/group/jepsen/docker.clj | 38 +++++++++++--- test/jepsen/src/group/jepsen/model.clj | 19 +++++-- test/jepsen/test/group/jepsen/model_test.clj | 52 ++++++++++++++++++- .../group/jepsen/retired_evidence_test.clj | 37 +++++++++++-- 9 files changed, 207 insertions(+), 25 deletions(-) create mode 100644 test/jepsen/decode_conflict_evidence.exs diff --git a/test/jepsen/README.md b/test/jepsen/README.md index 10a19d9..46348f3 100644 --- a/test/jepsen/README.md +++ b/test/jepsen/README.md @@ -77,7 +77,18 @@ The append-only conflict journal retains small operation records for the bounded campaign, not production ETS rows. Its path defaults to `/tmp/group-jepsen-conflict-evidence` inside each container and can be overridden with the driver's `:conflict_evidence_path` option. Corrupt or unreadable -journals fail closed. Checker qualification also loads the real Elixir +journals fail closed. At permanent retirement, the stopped-container collector +archives and decodes the journal alongside unexpected deaths. The checker +replays retired and surviving nodes' evidence together, without treating retired +owners as live. Conflict archives are bounded to 64 MiB with ten-second command +deadlines; missing or malformed archives fail the history. + +Each new history resets the running recorder through the harness socket after +DB restart and before workload mutations. This clears disk and in-memory +evidence together and initializes an empty journal, so repeated histories +cannot inherit earlier conflict coverage. + +Checker qualification also loads the real Elixir Owner/Driver harness, injects valid and forged death reasons, exercises a register interrupted before its reply, and checks the emitted EDN with the Clojure lifecycle oracle. diff --git a/test/jepsen/conflict_probe.exs b/test/jepsen/conflict_probe.exs index e15e2a4..f276ff6 100644 --- a/test/jepsen/conflict_probe.exs +++ b/test/jepsen/conflict_probe.exs @@ -13,7 +13,7 @@ defmodule Group.Jepsen.ConflictProbe do File.mkdir_p!(directory) try do - {winner, rejected, winner_events} = winner(directory) + {winner, rejected, winner_events, archive} = winner(directory) cases = [ {"historical winner subsequently unregistered and died", true, 1, "jepsen/registry/0", @@ -29,11 +29,13 @@ defmodule Group.Jepsen.ConflictProbe do scenarios = Enum.with_index(cases, fn {label, valid, revision, key, meta, pending}, index -> - events = victim(directory, index, revision, key, meta, pending) + {events, reset_events} = victim(directory, index, revision, key, meta, pending) %{ label: label, valid: valid, + archive: archive, + reset_evidence: reset_events, snapshots: %{ "n1" => %{conflict_evidence: events}, "n2" => %{conflict_evidence: winner_events} @@ -60,7 +62,7 @@ defmodule Group.Jepsen.ConflictProbe do end defp winner(directory) do - {group, evidence, driver, _path} = start(directory, "winner", "n2") + {group, evidence, driver, path} = start(directory, "winner", "n2") %{status: :ok, owner: %{token: token}} = mutate(driver, :register, 10) %{status: :fail, owner: %{token: rejected}} = @@ -70,10 +72,11 @@ defmodule Group.Jepsen.ConflictProbe do %{status: :ok} = GenServer.call(driver, {:kill, "owner"}) %{status: :ok} = GenServer.call(driver, {:kill, "rejected"}) events = ConflictEvidence.snapshot() + archive = File.read!(path) GenServer.stop(driver) GenServer.stop(evidence) Supervisor.stop(group) - {%{token: token, revision: 10}, %{token: rejected, revision: 100}, events} + {%{token: token, revision: 10}, %{token: rejected, revision: 100}, events, archive} end defp victim(directory, index, revision, key, winner, pending) do @@ -106,9 +109,27 @@ defmodule Group.Jepsen.ConflictProbe do # The exact registration/death obligations must survive a recorder restart. {:ok, evidence} = ConflictEvidence.start_link(conflict_evidence_path: path) ^events = ConflictEvidence.snapshot() + :ok = ConflictEvidence.reset() + [] = ConflictEvidence.snapshot() + "" = File.read!(path) + GenServer.stop(evidence) + {:ok, evidence} = ConflictEvidence.start_link(conflict_evidence_path: path) + [] = ConflictEvidence.snapshot() + + 1 = + ConflictEvidence.record(%{ + kind: :register, + token: "next-history", + key: 0, + cluster: nil, + revision: 1 + }) + + :ok = ConflictEvidence.reset() + reset_events = ConflictEvidence.snapshot() GenServer.stop(evidence) Supervisor.stop(group) - events + {events, reset_events} end defp wait(fun, remaining \\ 200) diff --git a/test/jepsen/decode_conflict_evidence.exs b/test/jepsen/decode_conflict_evidence.exs new file mode 100644 index 0000000..0fcc110 --- /dev/null +++ b/test/jepsen/decode_conflict_evidence.exs @@ -0,0 +1,6 @@ +System.put_env("GROUP_JEPSEN_LIBRARY_ONLY", "1") +Code.require_file("node.exs", __DIR__) + +[path] = System.argv() +events = path |> File.read!() |> Group.Jepsen.ConflictEvidence.decode() +IO.puts("CONFLICT-EVIDENCE " <> Group.Jepsen.EDN.encode(events)) diff --git a/test/jepsen/node.exs b/test/jepsen/node.exs index 152cc96..7ecc46d 100644 --- a/test/jepsen/node.exs +++ b/test/jepsen/node.exs @@ -434,6 +434,22 @@ defmodule Group.Jepsen.ConflictEvidence do def start_link(opts), do: GenServer.start_link(__MODULE__, opts, name: __MODULE__) def record(event), do: GenServer.call(__MODULE__, {:record, event}) def snapshot, do: GenServer.call(__MODULE__, :snapshot) + def reset, do: GenServer.call(__MODULE__, :reset) + + def decode(contents) do + if contents != "" and not String.ends_with?(contents, "\n"), + do: raise("truncated conflict evidence") + + contents + |> String.split("\n") + |> Enum.drop(-1) + |> Enum.with_index(1) + |> Enum.map(fn {line, sequence} -> + event = line |> Base.decode64!() |> :erlang.binary_to_term() + %{sequence: ^sequence} = event + event + end) + end @impl true def init(opts) do @@ -442,9 +458,7 @@ defmodule Group.Jepsen.ConflictEvidence do events = case File.read(path) do {:ok, contents} -> - contents - |> String.split("\n", trim: true) - |> Enum.map(&(&1 |> Base.decode64!() |> :erlang.binary_to_term())) + decode(contents) {:error, :enoent} -> [] @@ -467,6 +481,14 @@ defmodule Group.Jepsen.ConflictEvidence do def handle_call(:snapshot, _from, state) do {:reply, Enum.reverse(state.events), state} end + + # Called only during DB setup, after restart and before workload mutations. + # Removing the file externally would leave the restarted recorder's loaded + # evidence alive in memory and leak coverage into the next history. + def handle_call(:reset, _from, state) do + :ok = File.write(state.path, "", [:sync]) + {:reply, :ok, %{state | events: [], sequence: 0}} + end end defmodule Group.Jepsen.Owner do @@ -1444,6 +1466,10 @@ defmodule Group.Jepsen.Wire do ["ping"] -> %{status: :ok} + ["reset-conflict-evidence"] -> + :ok = Group.Jepsen.ConflictEvidence.reset() + %{status: :ok} + ["ready", expected] -> expected = String.to_integer(expected) diff --git a/test/jepsen/src/group/jepsen/db.clj b/test/jepsen/src/group/jepsen/db.clj index 702b41c..8898077 100644 --- a/test/jepsen/src/group/jepsen/db.clj +++ b/test/jepsen/src/group/jepsen/db.clj @@ -9,6 +9,10 @@ (docker/heal! (:nodes test)) (docker/restart! node) (docker/reset-oracle! node) + (group-client/wait-listening! node) + (let [response (group-client/request! node ["reset-conflict-evidence"])] + (when-not (= :ok (:status response)) + (throw (ex-info "conflict oracle reset failed" {:node node :response response})))) (group-client/wait-ready! node (count (:nodes test)))) (teardown! [_this test _node] diff --git a/test/jepsen/src/group/jepsen/docker.clj b/test/jepsen/src/group/jepsen/docker.clj index d0fc147..024d89f 100644 --- a/test/jepsen/src/group/jepsen/docker.clj +++ b/test/jepsen/src/group/jepsen/docker.clj @@ -1,5 +1,6 @@ (ns group.jepsen.docker - (:require [clojure.string :as str]) + (:require [clojure.edn :as edn] + [clojure.string :as str]) (:import (java.io File) (java.util.concurrent TimeUnit))) @@ -98,24 +99,45 @@ {:token token :reason reason})) (remove str/blank? (str/split-lines contents)))) -(defn retired-evidence! - "Reads the stopped container's durable oracle, without depending on its VM or socket. - Missing, truncated, unreadable, or oversized evidence is a qualification failure." - [node] +(defn decode-conflict-evidence! [file] + (let [output (shell! "sh" "-c" + (str "cd ../.. && exec env ERL_FLAGS='+S 2:2' mix run --no-start " + "test/jepsen/decode_conflict_evidence.exs \"$1\"") + "_" (.getAbsolutePath ^File file)) + prefix "CONFLICT-EVIDENCE " + line (first (filter #(str/starts-with? % prefix) (str/split-lines output)))] + (when-not line + (throw (ex-info "missing decoded conflict evidence" {}))) + (edn/read-string (subs line (count prefix))))) + +(defn collect-evidence! [node path bound decode] (let [file (File/createTempFile "group-jepsen-retired-" ".log")] (try (binding [*command-timeout-ms* 10000] (docker! "cp" - (str (container node) ":/tmp/group-jepsen-unexpected-deaths") + (str (container node) ":" path) (.getAbsolutePath file))) - (when (> (.length file) (* 8 1024 1024)) + (when (> (.length file) bound) (throw (ex-info "lifecycle evidence exceeds collection bound" {:node node}))) (let [contents (slurp file)] (when (and (seq contents) (not (str/ends-with? contents "\n"))) (throw (ex-info "truncated lifecycle evidence" {:node node}))) - {:node (name node) :unexpected-deaths (parse-unexpected-deaths contents)}) + (binding [*command-timeout-ms* 10000] + (decode file))) (finally (.delete file))))) +(defn retired-evidence! + "Reads the stopped container's durable oracle, without depending on its VM or socket. + Missing, truncated, unreadable, or oversized evidence is a qualification failure." + [node] + {:node (name node) + :unexpected-deaths + (collect-evidence! node "/tmp/group-jepsen-unexpected-deaths" (* 8 1024 1024) + #(parse-unexpected-deaths (slurp %))) + :conflict-evidence + (collect-evidence! node "/tmp/group-jepsen-conflict-evidence" (* 64 1024 1024) + decode-conflict-evidence!)}) + (defn ensure-firewall-chain! [node chain] (exec-sh! node diff --git a/test/jepsen/src/group/jepsen/model.clj b/test/jepsen/src/group/jepsen/model.clj index 222620a..4f31a8a 100644 --- a/test/jepsen/src/group/jepsen/model.clj +++ b/test/jepsen/src/group/jepsen/model.clj @@ -211,7 +211,22 @@ (when (not= expected-peers actual-peers) [node {:expected expected-peers, :actual actual-peers}])))) relevant-snapshots) - conflicts (conflict-analysis relevant-snapshots) + completions (remove history/invoke? history) + retirement-evidence (keep #(get-in % [:value :lifecycle-evidence]) completions) + ;; Retired nodes no longer contribute live owners or public views, but + ;; their durable claims and deaths remain obligations for this history. + conflict-snapshots + (into {} + (map (fn [[node evidence]] + [node {:conflict-evidence (->> evidence + (mapcat :conflict-evidence) + distinct + (sort-by :sequence) + vec)}])) + (group-by :node + (concat (map :value (successful-snapshots history)) + retirement-evidence))) + conflicts (conflict-analysis conflict-snapshots) transport-events (assoc (reduce #(merge-with + %1 %2) {} @@ -244,8 +259,6 @@ (not= 0 (:snapshot-staging-count internal))) [node internal])))) relevant-snapshots) - completions (remove history/invoke? history) - retirement-evidence (keep #(get-in % [:value :lifecycle-evidence]) completions) retired-nodes (set/difference (set (map name (:nodes test))) required-nodes) collected-nodes (set (map :node retirement-evidence)) evidence-errors (vec (keep #(get-in % [:value :evidence-error]) completions)) diff --git a/test/jepsen/test/group/jepsen/model_test.clj b/test/jepsen/test/group/jepsen/model_test.clj index 8d67335..056478f 100644 --- a/test/jepsen/test/group/jepsen/model_test.clj +++ b/test/jepsen/test/group/jepsen/model_test.clj @@ -3,6 +3,7 @@ [clojure.edn :as edn] [clojure.java.shell :as shell] [clojure.string :as str] + [group.jepsen.docker :as docker] [group.jepsen.model :as model])) (def test-map @@ -109,6 +110,30 @@ ;; Unique incarnation tokens break equal-revision ties without pid order. (valid (assoc-in conflict-journal [0 :revision] 2)))) +(deftest conflict-evidence-survives-permanent-retirement + (let [test (assoc test-map :terminal-nodes ["n2" "n3"] + :required-transport-events #{:registry-conflict-death}) + winner [{:kind :register :sequence 1 :token "winner" :cluster nil :key 0 :revision 10} + {:kind :result :sequence 2 :attempt 1 :status :ok}] + victim [{:kind :register :sequence 1 :token "victim" :cluster nil :key 0 :revision 1} + {:kind :death :sequence 2 :token "victim" :key "jepsen/registry/0" + :winner {:token "winner" :revision 10}}] + retired (fn [events] {:type :info :f :retire + :value {:lifecycle-evidence {:node "n1" :unexpected-deaths [] + :conflict-evidence events}}}) + terminal [(assoc-in (snapshot-op 1 "n2" ["n2" "n3"] [] (empty-registry) (empty-pg)) + [:value :conflict-evidence] victim) + (snapshot-op 2 "n3" ["n2" "n3"] [] (empty-registry) (empty-pg))] + healthy (model/analyze test (conj terminal (retired winner))) + bad (model/analyze (assoc test :required-transport-events #{}) + [(retired (assoc-in conflict-journal [3 :winner :revision] -100)) + (snapshot-op 2 "n2" ["n2" "n3"] [] (empty-registry) (empty-pg)) + (snapshot-op 3 "n3" ["n2" "n3"] [] (empty-registry) (empty-pg))])] + (is (:valid? healthy)) + (is (= 1 (get-in healthy [:transport-events :registry-conflict-death]))) + (is (false? (:valid? bad))) + (is (= 1 (count (:invalid-conflict-deaths bad)))))) + (deftest conflict-statistics-alone-do-not-discharge-deaths (let [history [(assoc-in (snapshot-op 1 "n1" [] (empty-registry) (empty-pg)) [:value :transport-events] {:registry-conflict-death 99}) @@ -129,14 +154,37 @@ (str/split-lines out))) scenarios (when line (edn/read-string (subs line 15)))] (is (some? scenarios) (str out err)) - (doseq [{:keys [label valid snapshots]} scenarios] + (doseq [{:keys [label valid snapshots reset-evidence]} scenarios] (let [history (mapv (fn [index node] (assoc-in (snapshot-op index node [] (empty-registry) (empty-pg)) [:value :conflict-evidence] (get-in snapshots [node :conflict-evidence]))) (range 3) ["n1" "n2" "n3"]) result (model/analyze test-map history)] - (is (= valid (:valid? result)) (str label ": " result)))))) + (is (= valid (:valid? result)) (str label ": " result))) + (let [empty-history (mapv #(assoc-in (snapshot-op %1 %2 [] (empty-registry) (empty-pg)) + [:value :conflict-evidence] reset-evidence) + (range 3) ["n1" "n2" "n3"]) + result (model/analyze (assoc test-map :required-transport-events #{:registry-conflict-death}) + empty-history)] + (is (false? (:valid? result)) "reset history cannot reuse conflict coverage"))) + (when scenarios + (with-redefs [docker/docker! + (fn [& args] + (spit (last args) + (if (str/includes? (second args) "conflict-evidence") + (:archive (first scenarios)) "")))] + (let [archive (docker/retired-evidence! "n2") + test (assoc test-map :terminal-nodes ["n1" "n3"]) + scenario (first scenarios) + history [{:type :info :f :retire :value {:lifecycle-evidence archive}} + (assoc-in (snapshot-op 1 "n1" ["n1" "n3"] [] (empty-registry) (empty-pg)) + [:value :conflict-evidence] + (get-in scenario [:snapshots "n1" :conflict-evidence])) + (snapshot-op 2 "n3" ["n1" "n3"] [] (empty-registry) (empty-pg))]] + (is (= (get-in scenario [:snapshots "n2" :conflict-evidence]) + (:conflict-evidence archive))) + (is (:valid? (model/analyze test history)))))))) (deftest accepts-an-exact-converged-multi-cluster-view (let [owners [(owner "a" [(registration nil 0 1) (registration "red" 1 2)] []) diff --git a/test/jepsen/test/group/jepsen/retired_evidence_test.clj b/test/jepsen/test/group/jepsen/retired_evidence_test.clj index 3491e73..5f738c5 100644 --- a/test/jepsen/test/group/jepsen/retired_evidence_test.clj +++ b/test/jepsen/test/group/jepsen/retired_evidence_test.clj @@ -1,6 +1,8 @@ (ns group.jepsen.retired-evidence-test (:require [clojure.test :refer :all] [group.jepsen.docker :as docker] + [group.jepsen.client :as client] + [group.jepsen.db :as group-db] [group.jepsen.nemesis :as group-nemesis] [jepsen.db :as db] [jepsen.nemesis :as nemesis])) @@ -11,15 +13,18 @@ (fn [& args] (swap! calls conj args) (is (= 10000 docker/*command-timeout-ms*)) - (spit (last args) "n1/boot/owner/1\t:boom\n"))] - (is (= {:node "n1" :unexpected-deaths [{:token "n1/boot/owner/1" :reason ":boom"}]} + (spit (last args) "n1/boot/owner/1\t:boom\n")) + docker/decode-conflict-evidence! (constantly [])] + (is (= {:node "n1" :unexpected-deaths [{:token "n1/boot/owner/1" :reason ":boom"}] + :conflict-evidence []} (docker/retired-evidence! "n1"))) (is (= ["cp" "group-jepsen-n1:/tmp/group-jepsen-unexpected-deaths"] (vec (take 2 (first @calls)))))))) (deftest collector-distinguishes-empty-evidence-from-loss (doseq [contents ["" "bad\n" "token\t:boom"]] - (with-redefs [docker/docker! (fn [& args] (spit (last args) contents))] + (with-redefs [docker/docker! (fn [& args] (spit (last args) contents)) + docker/decode-conflict-evidence! (constantly [])] (if (= "" contents) (is (= [] (:unexpected-deaths (docker/retired-evidence! "n1")))) (is (thrown? Exception (docker/retired-evidence! "n1")))))) @@ -33,6 +38,32 @@ (is (re-find #": > /tmp/group-jepsen-unexpected-deaths" script)))] (docker/reset-oracle! "n1"))) +(deftest workload-reset-clears-the-running-conflict-recorder + (let [calls (atom [])] + (with-redefs [docker/heal! (fn [_]) + docker/restart! #(swap! calls conj [:restart %]) + docker/reset-oracle! #(swap! calls conj [:disk-reset %]) + client/wait-listening! #(swap! calls conj [:listening %]) + client/request! (fn [node fields] + (swap! calls conj [node fields]) + {:status :ok}) + client/wait-ready! (fn [node _] (swap! calls conj [:ready node]))] + (dotimes [_ 2] + (db/setup! (group-db/db) {:nodes ["n1"]} "n1")) + (is (= (vec (mapcat identity (repeat 2 [[:restart "n1"] [:disk-reset "n1"] + [:listening "n1"] + ["n1" ["reset-conflict-evidence"]] + [:ready "n1"]]))) + @calls))))) + +(deftest conflict-archive-decode-fails-closed + (doseq [contents ["not-base64\n" "truncated"]] + (with-redefs [docker/docker! + (fn [& args] + (spit (last args) + (if (.contains (second args) "conflict-evidence") contents "")))] + (is (thrown? Exception (docker/retired-evidence! "n1")))))) + (deftest retirement-captures-after-stop-even-if-node-was-already-unavailable (doseq [running? [true false]] (let [calls (atom []) From c6299b260ef6adf7944de650ca976b44a3003f8d Mon Sep 17 00:00:00 2001 From: Jason Stiebs Date: Fri, 11 Sep 2026 16:48:10 -0500 Subject: [PATCH 8/8] Run live checker qualification with Elixir --- test/jepsen/README.md | 11 +- test/jepsen/qualification-result-test.sh | 36 ----- test/jepsen/qualification-result.sh | 19 --- test/jepsen/qualify.exs | 99 +++++++++++++ test/jepsen/qualify.sh | 50 +------ test/jepsen_qualification_cache_test.exs | 1 + test/jepsen_qualification_runner_test.exs | 170 ++++++++++++++++++++++ 7 files changed, 281 insertions(+), 105 deletions(-) delete mode 100755 test/jepsen/qualification-result-test.sh delete mode 100644 test/jepsen/qualification-result.sh create mode 100644 test/jepsen/qualify.exs create mode 100644 test/jepsen_qualification_runner_test.exs diff --git a/test/jepsen/README.md b/test/jepsen/README.md index 199c9e7..5f43195 100644 --- a/test/jepsen/README.md +++ b/test/jepsen/README.md @@ -109,9 +109,18 @@ count, concurrency, keys, owners, and recovery time. Run mutation qualification plus live positive- and negative-checker tests: ```bash -test/jepsen/qualify.sh +elixir test/jepsen/qualify.exs ``` +`qualify.sh` remains a thin compatibility launcher for the same Elixir script. +The script owns artifact creation, command sequencing, and result validation. +Each live run still uses GNU `timeout` with a three-minute deadline and a +30-second TERM/KILL grace period, and streams its output to a separate log. +Qualification requires both the expected exit status and a fresh checker record +confirming that the intended corruption was actually detected; a failed command +alone never counts. Executable ExUnit regressions exercise this runner with +stubbed external commands as part of normal `mix test`, without Docker. + This runs every mutation defined by `test/mutation/run.exs`, then verifies that a healthy live history is accepted and deliberately injected owner-death, internal-index, stranded snapshot-cursor, registry claim/projection, and diff --git a/test/jepsen/qualification-result-test.sh b/test/jepsen/qualification-result-test.sh deleted file mode 100755 index 3dda51b..0000000 --- a/test/jepsen/qualification-result-test.sh +++ /dev/null @@ -1,36 +0,0 @@ -#!/usr/bin/env bash -set -euo pipefail -script_dir="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" -source "${script_dir}/qualification-result.sh" -work="$(mktemp -d "${script_dir}/qualification-test.XXXXXX")" -trap 'rm -rf "${work}"' EXIT - -probe() { - local want="$1" expectation="$2" mode="$3" status="$4" valid="$5" qualified="$6" - local result="${work}/result" - rm -f "${result}" - # A bounded stand-in for the CLI: no Docker, JVM or network required. - if [[ "${valid}" != absent ]]; then - printf 'group-qualification-v1\t%s\t%s\t%s\n' "${mode}" "${valid}" "${qualified}" >"${result}" - fi - local actual=reject - if qualification_result "${expectation}" "${mode}" "${status}" "${result}"; then - actual=accept - fi - [[ "${actual}" == "${want}" ]] || { echo "unexpected classification"; exit 1; } -} - -probe accept pass none 0 true true -for mode in unexpected-death internal-index cursor-marker registry-projection terminal-unavailable; do - probe accept fail "${mode}" 1 false true - probe reject fail "${mode}" 1 false false -done -probe reject pass none 1 false false -probe reject pass none 1 absent false -probe reject fail internal-index 0 true true -for status in 1 124 125 137 127; do - probe reject fail internal-index "${status}" absent false -done -probe reject fail internal-index 124 false true -probe reject fail internal-index 125 false true -echo "qualification status regression passed" diff --git a/test/jepsen/qualification-result.sh b/test/jepsen/qualification-result.sh deleted file mode 100644 index e4cdf21..0000000 --- a/test/jepsen/qualification-result.sh +++ /dev/null @@ -1,19 +0,0 @@ -#!/usr/bin/env bash - -# Only a normal CLI exit and a fresh, completed checker decision qualify. -# In particular timeout (124/137), timeout invocation (125), and JVM/Docker -# failures are not negative checker evidence, even if an artifact exists. -qualification_result() { - local expectation="$1" corruption="$2" status="$3" result="$4" - local expected_status=1 valid=false - if [[ "${expectation}" == pass ]]; then - expected_status=0 - valid=true - fi - if [[ "${status}" -ne "${expected_status}" ]] || - [[ ! -f "${result}" ]] || - [[ "$(cat "${result}")" != "$(printf 'group-qualification-v1\t%s\t%s\ttrue' "${corruption}" "${valid}")" ]]; then - echo "Jepsen qualification failed: ${expectation}/${corruption}, exit ${status}; missing, mismatched, or unqualified checker result" >&2 - return 1 - fi -} diff --git a/test/jepsen/qualify.exs b/test/jepsen/qualify.exs new file mode 100644 index 0000000..14da196 --- /dev/null +++ b/test/jepsen/qualify.exs @@ -0,0 +1,99 @@ +defmodule Group.Jepsen.Qualification do + @moduledoc false + + # An intentionally corrupted history must be invalid, but still demonstrate + # that its intended corruption was injected and detected by the checker. + @checks [ + {"none", true}, + {"unexpected-death", false}, + {"internal-index", false}, + {"cursor-marker", false}, + {"registry-projection", false}, + {"terminal-unavailable", false} + ] + + def run do + repo = Path.expand("../..", __DIR__) + cache = Path.join(__DIR__, ".cache") + File.mkdir_p!(cache) + suffix = Base.url_encode64(:crypto.strong_rand_bytes(12), padding: false) + artifacts = Path.join(cache, "qualification.#{suffix}") + File.mkdir!(artifacts) + File.chmod!(artifacts, 0o700) + IO.puts("qualification artifacts: #{artifacts}") + + {_output, status} = + System.cmd("mix", ["run", "test/mutation/run.exs"], + cd: repo, + into: IO.stream(), + stderr_to_stdout: true + ) + + if status != 0, do: System.halt(status) + + Enum.each(@checks, fn {corruption, expected_valid?} -> + qualify!(repo, artifacts, corruption, expected_valid?) + end) + + IO.puts("mutation and live checker qualification passed") + end + + defp qualify!(repo, artifacts, corruption, expected_valid?) do + log = Path.join(artifacts, "#{corruption}.log") + result = Path.join(artifacts, "#{corruption}.result") + + # Keep the existing process-tree deadline and TERM/KILL grace period. + # Streaming output to disk avoids retaining a live history on this VM's heap. + {_output, status} = + System.cmd( + "timeout", + [ + "--signal=TERM", + "--kill-after=30", + "180", + Path.join(__DIR__, "run.sh"), + "test", + "--no-ssh", + "--nodes", + "n1,n2,n3", + "--concurrency", + "2n", + "--time-limit", + "6", + "--fault-interval", + "1", + "--recovery-time", + "5", + "--transport", + "distribution", + "--scenario", + "mixed", + "--corruption", + corruption + ], + cd: repo, + env: [ + {"GROUP_JEPSEN_SKIP_CHECKER", "1"}, + {"GROUP_JEPSEN_QUALIFICATION_RESULT", result} + ], + into: File.stream!(log, [:write, :binary]), + stderr_to_stdout: true + ) + + expected_status = if expected_valid?, do: 0, else: 1 + + unless status == expected_status do + raise "#{corruption}: expected exit #{expected_status}, got #{status}; see #{log}" + end + + expected_record = "group-qualification-v1\t#{corruption}\t#{expected_valid?}\ttrue\n" + + unless File.read!(result) == expected_record do + raise "#{corruption}: missing or mismatched checker evidence; see #{result} and #{log}" + end + + IO.puts("qualified #{corruption} (#{log})") + end +end + +Group.Jepsen.Qualification.run() diff --git a/test/jepsen/qualify.sh b/test/jepsen/qualify.sh index 78f26d0..ebc3893 100755 --- a/test/jepsen/qualify.sh +++ b/test/jepsen/qualify.sh @@ -2,52 +2,4 @@ set -euo pipefail script_dir="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" -repo_dir="$(cd "${script_dir}/../.." && pwd)" -mkdir -p "${script_dir}/.cache" -artifact_dir="$(mktemp -d "${script_dir}/.cache/qualification.XXXXXX")" -source "${script_dir}/qualification-result.sh" - -cd "${repo_dir}" - -mix run test/mutation/run.exs - -export GROUP_JEPSEN_SKIP_CHECKER=1 - -run_jepsen() { - local expectation="$1" - local corruption="$2" - local log="${artifact_dir}/${expectation}-${corruption}.log" - local result="${artifact_dir}/${expectation}-${corruption}.result" - local status=0 - - set +e - GROUP_JEPSEN_QUALIFICATION_RESULT="${result}" \ - timeout --signal=TERM --kill-after=30 180 "${script_dir}/run.sh" test \ - --no-ssh \ - --nodes n1,n2,n3 \ - --concurrency 2n \ - --time-limit 6 \ - --fault-interval 1 \ - --recovery-time 5 \ - --transport distribution \ - --scenario mixed \ - --corruption "${corruption}" >"${log}" 2>&1 - status=$? - set -e - - if ! qualification_result "${expectation}" "${corruption}" "${status}" "${result}"; then - echo "see ${log} and ${result}" >&2 - return 1 - fi - - echo "${expectation}: ${corruption} (${log})" -} - -run_jepsen pass none -run_jepsen fail unexpected-death -run_jepsen fail internal-index -run_jepsen fail cursor-marker -run_jepsen fail registry-projection -run_jepsen fail terminal-unavailable - -echo "mutation and live checker qualification passed" +exec elixir "${script_dir}/qualify.exs" "$@" diff --git a/test/jepsen_qualification_cache_test.exs b/test/jepsen_qualification_cache_test.exs index 62de269..d5ce0d7 100644 --- a/test/jepsen_qualification_cache_test.exs +++ b/test/jepsen_qualification_cache_test.exs @@ -15,6 +15,7 @@ defmodule Group.JepsenQualificationCacheTest do File.mkdir_p!(script_dir) File.mkdir_p!(bin) File.cp!("test/jepsen/qualify.sh", Path.join(script_dir, "qualify.sh")) + File.cp!("test/jepsen/qualify.exs", Path.join(script_dir, "qualify.exs")) if existing_cache? do File.mkdir_p!(cache) diff --git a/test/jepsen_qualification_runner_test.exs b/test/jepsen_qualification_runner_test.exs new file mode 100644 index 0000000..e165529 --- /dev/null +++ b/test/jepsen_qualification_runner_test.exs @@ -0,0 +1,170 @@ +defmodule Group.JepsenQualificationRunnerTest do + use ExUnit.Case, async: true + + @moduletag :local + @moduletag tmp_dir: System.pid() + + @corruptions ~w(none unexpected-death internal-index cursor-marker registry-projection terminal-unavailable) + + setup %{tmp_dir: directory} do + repo = Path.expand("checkout with spaces", directory) + scripts = Path.join(repo, "test/jepsen") + bin = Path.expand("bin", directory) + probes = Path.expand("probes", directory) + + for path <- [scripts, bin, probes], do: File.mkdir_p!(path) + script = Path.join(scripts, "qualify.exs") + File.cp!("test/jepsen/qualify.exs", script) + + # Exercise the actual executable runner, replacing only the external + # mutation/Jepsen commands. Neither Docker nor the real mutation campaign runs. + for command <- ["mix", "timeout"] do + path = Path.join(bin, command) + File.write!(path, command_stub()) + File.chmod!(path, 0o755) + end + + {:ok, repo: repo, scripts: scripts, script: script, bin: bin, probes: probes} + end + + test "qualifies the healthy baseline and all five corruptions with bounded commands", context do + {output, status} = run_qualification(context) + assert status == 0, output + assert calls(context) == ["mutations" | @corruptions] + assert [artifacts] = artifact_directories(context) + + for mode <- @corruptions do + [cwd, skip_checker, result | args] = invocation(context, mode) + assert cwd == context.repo + assert skip_checker == "1" + assert result == Path.join(artifacts, "#{mode}.result") + + assert args == [ + "--signal=TERM", + "--kill-after=30", + "180", + Path.join(context.scripts, "run.sh"), + "test", + "--no-ssh", + "--nodes", + "n1,n2,n3", + "--concurrency", + "2n", + "--time-limit", + "6", + "--fault-interval", + "1", + "--recovery-time", + "5", + "--transport", + "distribution", + "--scenario", + "mixed", + "--corruption", + mode + ] + + assert File.read!(Path.join(artifacts, "#{mode}.log")) == "probe log #{mode}\n" + valid? = mode == "none" + assert File.read!(result) == "group-qualification-v1\t#{mode}\t#{valid?}\ttrue\n" + end + end + + test "a failed mutation baseline propagates without running live qualification", context do + assert {_, 42} = run_qualification(context, baseline_status: 42) + assert calls(context) == ["mutations"] + end + + test "a failing healthy history stops before any corruption", context do + assert {_, 1} = run_qualification(context, mode: "none", status: 1) + assert calls(context) == ["mutations", "none"] + end + + for {label, status, record} <- [ + {"unexpected acceptance", 0, "wrong-validity"}, + {"unqualified rejection", 1, "unqualified"}, + {"wrong corruption", 1, "wrong-mode"}, + {"wrong schema", 1, "wrong-version"}, + {"malformed record", 1, "malformed"}, + {"missing result", 1, "missing"}, + {"timeout with completed evidence", 124, "valid"}, + {"timeout invocation failure", 125, "valid"}, + {"killed process", 137, "missing"}, + {"missing executable", 127, "missing"} + ] do + @tag status: status, record: record + test "rejects #{label} and does not proceed to the next corruption", context do + assert {_, 1} = run_qualification(context, status: context.status, record: context.record) + assert calls(context) == ["mutations", "none", "unexpected-death"] + end + end + + test "a prior successful run cannot supply a missing result for the next run", context do + assert {_, 0} = run_qualification(context) + [first] = artifact_directories(context) + assert {_, 1} = run_qualification(context, record: "missing") + assert length(artifact_directories(context)) == 2 + [_, _, result | _] = invocation(context, "unexpected-death") + refute Path.dirname(result) == first + refute File.exists?(result) + assert File.exists?(Path.join(first, "unexpected-death.result")) + end + + defp run_qualification(context, opts \\ []) do + System.cmd("elixir", [context.script], + env: [ + {"PATH", context.bin <> ":" <> System.fetch_env!("PATH")}, + {"ERL_FLAGS", "+S 1:1"}, + {"QUALIFICATION_PROBE_DIR", context.probes}, + {"QUALIFICATION_PROBE_MODE", Keyword.get(opts, :mode, "unexpected-death")}, + {"QUALIFICATION_PROBE_RECORD", Keyword.get(opts, :record, "valid")}, + {"QUALIFICATION_PROBE_STATUS", opts[:status] && to_string(opts[:status])}, + {"QUALIFICATION_BASELINE_STATUS", to_string(Keyword.get(opts, :baseline_status, 0))} + ], + stderr_to_stdout: true + ) + end + + defp calls(context) do + context.probes |> Path.join("calls") |> File.read!() |> String.split("\n", trim: true) + end + + defp invocation(context, mode) do + context.probes |> Path.join("#{mode}.txt") |> File.read!() |> String.split("\n") + end + + defp artifact_directories(context) do + Path.wildcard(Path.join(context.scripts, ".cache/qualification.*")) + end + + defp command_stub do + ~S""" + #!/usr/bin/env elixir + args = System.argv() + probes = System.fetch_env!("QUALIFICATION_PROBE_DIR") + mode = if args == ["run", "test/mutation/run.exs"], do: "mutations", else: List.last(args) + File.write!(Path.join(probes, "calls"), mode <> "\n", [:append]) + result = System.get_env("GROUP_JEPSEN_QUALIFICATION_RESULT", "absent") + skip = System.get_env("GROUP_JEPSEN_SKIP_CHECKER", "unset") + File.write!(Path.join(probes, "#{mode}.txt"), Enum.join([File.cwd!(), skip, result | args], "\n")) + IO.puts("probe log #{mode}") + + if mode == "mutations" do + System.halt(String.to_integer(System.fetch_env!("QUALIFICATION_BASELINE_STATUS"))) + end + + targeted? = mode == System.fetch_env!("QUALIFICATION_PROBE_MODE") + status = if mode == "none", do: "0", else: "1" + status = if targeted?, do: System.get_env("QUALIFICATION_PROBE_STATUS") || status, else: status + kind = if targeted?, do: System.fetch_env!("QUALIFICATION_PROBE_RECORD"), else: "valid" + valid? = mode == "none" + valid? = if kind == "wrong-validity", do: not valid?, else: valid? + record_mode = if kind == "wrong-mode", do: "different-corruption", else: mode + version = if kind == "wrong-version", do: "group-qualification-v2", else: "group-qualification-v1" + record = "#{version}\t#{record_mode}\t#{valid?}\t#{kind != "unqualified"}\n" + record = if kind == "malformed", do: "not a qualification record\n", else: record + unless kind == "missing", do: File.write!(result, record) + System.halt(String.to_integer(status)) + """ + end +end