From 9c2a7fbe57e15689ecdb2c383b9de1d0cb33531a Mon Sep 17 00:00:00 2001 From: Jason Stiebs Date: Fri, 11 Sep 2026 12:33:01 -0500 Subject: [PATCH 1/2] 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 7b88453b79bd13726881cfd10c3d2bd3ba8e1326 Mon Sep 17 00:00:00 2001 From: Jason Stiebs Date: Fri, 11 Sep 2026 16:13:53 -0500 Subject: [PATCH 2/2] 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 [])