Skip to content

Commit edccc34

Browse files
authored
Merge pull request #4422 from Yue021130/codex/3700-claim-work-parity
test(coordination): characterize file-backed claim_work with a parity fixture
2 parents e993fb5 + c7a772f commit edccc34

1 file changed

Lines changed: 361 additions & 0 deletions

File tree

Lines changed: 361 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,361 @@
1+
"""Characterize the shipped file-backed ``claim_work`` path against the contract.
2+
3+
``tests/control_plane/test_coordination_executor.py`` proves the authority
4+
semantics with an in-memory provider, and
5+
``test_coordination_file_provider.py`` proves the storage verbs. This fixture
6+
closes the gap between them: one scenario matrix drives the real
7+
``CoordinationAuthorityExecutor`` through the shipped
8+
``FileCoordinationProvider``.
9+
10+
The fixture is provider-neutral by construction. Every scenario speaks only the
11+
storage protocol (``store_identity`` / ``load`` / ``compare_and_put``) and
12+
receives providers from a handle factory, so registering another provider later
13+
requires no change to the matrix or to the expectations.
14+
15+
Expected outcomes are declared from the coordination domain contract (RFC
16+
section 10, checks 1-5) before anything runs. They are never read back from
17+
provider output, and no authority rule is re-derived here: the executor is the
18+
only decision maker under test.
19+
20+
In scope: characterization only. No production code changes, no new provider,
21+
no default-provider or public-behavior change, and no live NoKV, credential,
22+
service-startup, or provider-promotion surface.
23+
"""
24+
25+
from __future__ import annotations
26+
27+
import copy
28+
from dataclasses import dataclass, field
29+
from typing import Any, Callable
30+
31+
import pytest
32+
33+
from loopx.control_plane.coordination.executor import (
34+
CoordinationAuthorityExecutor,
35+
sample_claim_envelope,
36+
)
37+
from loopx.control_plane.coordination.file_provider import FileCoordinationProvider
38+
from loopx.control_plane.coordination.head import bootstrap_head
39+
40+
41+
GOAL_ID = "goal-a"
42+
# Fixed wall clock: every lease expiry in this matrix is minted from it.
43+
NOW = 1_800_000_000.0
44+
ELIGIBLE_AGENTS = ("agent-a", "agent-b")
45+
46+
47+
# ---- synthetic, public-safe head fixtures -----------------------------------
48+
49+
50+
def eligibility(allowed: tuple[str, ...] = ELIGIBLE_AGENTS) -> dict[str, Any]:
51+
return {
52+
"authorization_projection_revision": 3,
53+
"authorization_projection_digest": "sha256:bootstrap-auth",
54+
"allowed_agent_ids": list(allowed),
55+
"dependencies_satisfied": True,
56+
"dependency_revision": 12,
57+
"gates_open": True,
58+
"gate_revision": 5,
59+
}
60+
61+
62+
def todo(**overrides: Any) -> dict[str, Any]:
63+
base: dict[str, Any] = {
64+
"todo_revision": 7,
65+
"status": "open",
66+
"claimed_by": None,
67+
"eligibility": eligibility(),
68+
"repository": "git:example/repo",
69+
"code_revision": "0123456789abcdef",
70+
"last_lease_epoch": 6,
71+
}
72+
base.update(overrides)
73+
return base
74+
75+
76+
def claim(
77+
agent: str = "agent-a",
78+
todo_id: str = "todo-1",
79+
operation_id: str | None = None,
80+
**overrides: Any,
81+
) -> dict[str, Any]:
82+
return sample_claim_envelope(
83+
goal_id=GOAL_ID,
84+
operation_id=operation_id or f"op-{agent}-{todo_id}",
85+
agent_id=agent,
86+
device_id=f"dev-{agent}",
87+
todo_id=todo_id,
88+
expected_todo_revision=7,
89+
expected_preconditions={
90+
"authorization_projection_revision": 3,
91+
"authorization_projection_digest": "sha256:bootstrap-auth",
92+
"dependency_revision": 12,
93+
"gate_revision": 5,
94+
},
95+
lease_ttl_seconds=600,
96+
**overrides,
97+
)
98+
99+
100+
def bootstrap(provider: Any, todo_ids: tuple[str, ...] = ("todo-1", "todo-2")) -> None:
101+
head = bootstrap_head(
102+
GOAL_ID,
103+
{todo_id: todo() for todo_id in todo_ids},
104+
store_binding=provider.store_identity(),
105+
)
106+
assert provider.compare_and_put(0, head)["result"] == "applied"
107+
108+
109+
def executor_for(provider: Any) -> CoordinationAuthorityExecutor:
110+
return CoordinationAuthorityExecutor(provider, goal_id=GOAL_ID, now=lambda: NOW)
111+
112+
113+
class RecordingProvider:
114+
"""Storage-only recorder: delegates every verb and records CAS expectations.
115+
116+
It holds no authority semantics; it exists so a scenario can prove that a
117+
stale generation reached ``compare_and_put`` without duplicating a write.
118+
"""
119+
120+
def __init__(self, inner: Any) -> None:
121+
self._inner = inner
122+
self.cas_expectations: list[int] = []
123+
124+
def store_identity(self) -> str:
125+
return self._inner.store_identity()
126+
127+
def load(self) -> tuple[dict[str, Any] | None, int]:
128+
return self._inner.load()
129+
130+
def compare_and_put(
131+
self, expected_provider_generation: int, head: dict[str, Any]
132+
) -> dict[str, Any]:
133+
self.cas_expectations.append(expected_provider_generation)
134+
return self._inner.compare_and_put(expected_provider_generation, head)
135+
136+
137+
# ---- observation and expectation --------------------------------------------
138+
139+
140+
@dataclass
141+
class Observation:
142+
"""What one scenario actually produced, in contract terms only."""
143+
144+
outcomes: list[dict[str, Any]] = field(default_factory=list)
145+
head: dict[str, Any] = field(default_factory=dict)
146+
flags: dict[str, bool] = field(default_factory=dict)
147+
148+
149+
@dataclass(frozen=True)
150+
class Expectation:
151+
"""The contract-derived verdict one scenario must produce."""
152+
153+
results: tuple[str, ...]
154+
reasons: tuple[str | None, ...]
155+
authority_revision: int
156+
receipts: tuple[str, ...]
157+
claimed_by: tuple[tuple[str, str | None], ...]
158+
todo_revisions: tuple[tuple[str, int], ...]
159+
flags: tuple[tuple[str, bool], ...] = ()
160+
161+
162+
# ---- scenarios --------------------------------------------------------------
163+
164+
165+
def same_target_competition(handles: Callable[[], Any]) -> Observation:
166+
"""Check 1: two actors on one target leave exactly one winner."""
167+
168+
provider = handles()
169+
bootstrap(provider)
170+
executor = executor_for(provider)
171+
observation = Observation()
172+
observation.outcomes.append(executor.apply(claim("agent-a", "todo-1")))
173+
observation.outcomes.append(executor.apply(claim("agent-b", "todo-1")))
174+
observation.head, _ = provider.load()
175+
return observation
176+
177+
178+
def independent_targets_rebase(handles: Callable[[], Any]) -> Observation:
179+
"""Check 2: independent targets both apply through internal rebase."""
180+
181+
provider = handles()
182+
bootstrap(provider)
183+
observation = Observation()
184+
observation.outcomes.append(
185+
executor_for(provider).apply(claim("agent-a", "todo-1"))
186+
)
187+
observation.outcomes.append(
188+
executor_for(provider).apply(claim("agent-b", "todo-2"))
189+
)
190+
observation.head, _ = provider.load()
191+
return observation
192+
193+
194+
def replay_after_interleaved_write(handles: Callable[[], Any]) -> Observation:
195+
"""Check 3: `A -> B -> replay A` returns the exact original receipt.
196+
197+
The replay runs through a freshly constructed executor and a fresh provider
198+
handle, so nothing is served from in-process state.
199+
"""
200+
201+
provider = handles()
202+
bootstrap(provider)
203+
executor = executor_for(provider)
204+
observation = Observation()
205+
request_a = claim("agent-a", "todo-1", operation_id="op-A")
206+
first = executor.apply(request_a)
207+
observation.outcomes.append(first)
208+
observation.outcomes.append(
209+
executor.apply(claim("agent-b", "todo-2", operation_id="op-B"))
210+
)
211+
reconstructed = executor_for(handles())
212+
replay = reconstructed.apply(copy.deepcopy(request_a))
213+
observation.outcomes.append(replay)
214+
observation.flags["replay_returns_original_receipt"] = (
215+
replay.get("original_receipt") == first.get("original_receipt")
216+
and replay.get("original_receipt") is not None
217+
)
218+
observation.head, _ = provider.load()
219+
return observation
220+
221+
222+
def operation_identity_reuse(handles: Callable[[], Any]) -> Observation:
223+
"""Check 5: one operation id with different semantics changes nothing."""
224+
225+
provider = handles()
226+
bootstrap(provider)
227+
executor = executor_for(provider)
228+
observation = Observation()
229+
request = claim("agent-a", "todo-1", operation_id="op-A")
230+
observation.outcomes.append(executor.apply(request))
231+
before, _ = provider.load()
232+
mutated = copy.deepcopy(request)
233+
mutated["command"]["lease_ttl_seconds"] = 601
234+
observation.outcomes.append(executor.apply(mutated))
235+
observation.head, _ = provider.load()
236+
observation.flags["state_unchanged_after_rejection"] = observation.head == before
237+
return observation
238+
239+
240+
def stale_generation_does_not_duplicate(handles: Callable[[], Any]) -> Observation:
241+
"""Stale provider generation is refused rather than replayed.
242+
243+
A handle observes a generation, the document advances behind it, and the
244+
stale expectation is then offered for a raw CAS. The refused attempt must
245+
not write, and the executor's bounded reload must land exactly one further
246+
authority transition with no duplicate receipt.
247+
"""
248+
249+
provider = handles()
250+
bootstrap(provider)
251+
stale = RecordingProvider(handles())
252+
stale_generation = stale.load()[1]
253+
254+
observation = Observation()
255+
observation.outcomes.append(
256+
executor_for(provider).apply(claim("agent-a", "todo-1"))
257+
)
258+
current_head, current_generation = provider.load()
259+
stale_attempt = stale.compare_and_put(stale_generation, current_head)
260+
observation.flags["stale_cas_refused"] = stale_attempt["result"] == "conflict"
261+
observation.flags["stale_cas_wrote_nothing"] = provider.load()[1] == current_generation
262+
observation.flags["stale_expectation_offered"] = bool(
263+
stale.cas_expectations and stale.cas_expectations[0] == stale_generation
264+
)
265+
266+
observation.outcomes.append(
267+
executor_for(stale).apply(claim("agent-b", "todo-2"))
268+
)
269+
observation.head, _ = provider.load()
270+
return observation
271+
272+
273+
SCENARIOS: dict[str, Callable[[Callable[[], Any]], Observation]] = {
274+
"same_target_competition": same_target_competition,
275+
"independent_targets_rebase": independent_targets_rebase,
276+
"replay_after_interleaved_write": replay_after_interleaved_write,
277+
"operation_identity_reuse": operation_identity_reuse,
278+
"stale_generation_does_not_duplicate": stale_generation_does_not_duplicate,
279+
}
280+
281+
EXPECTED: dict[str, Expectation] = {
282+
"same_target_competition": Expectation(
283+
results=("applied", "conflict"),
284+
reasons=(None, "todo_revision_mismatch"),
285+
authority_revision=1,
286+
receipts=("op-agent-a-todo-1",),
287+
claimed_by=(("todo-1", "agent-a"), ("todo-2", None)),
288+
todo_revisions=(("todo-1", 8), ("todo-2", 7)),
289+
),
290+
"independent_targets_rebase": Expectation(
291+
results=("applied", "applied"),
292+
reasons=(None, None),
293+
authority_revision=2,
294+
receipts=("op-agent-a-todo-1", "op-agent-b-todo-2"),
295+
claimed_by=(("todo-1", "agent-a"), ("todo-2", "agent-b")),
296+
todo_revisions=(("todo-1", 8), ("todo-2", 8)),
297+
),
298+
"replay_after_interleaved_write": Expectation(
299+
results=("applied", "applied", "already_applied"),
300+
reasons=(None, None, None),
301+
authority_revision=2,
302+
receipts=("op-A", "op-B"),
303+
claimed_by=(("todo-1", "agent-a"), ("todo-2", "agent-b")),
304+
todo_revisions=(("todo-1", 8), ("todo-2", 8)),
305+
# A replayed operation is a read, not a second transition.
306+
flags=(("replay_returns_original_receipt", True),),
307+
),
308+
"operation_identity_reuse": Expectation(
309+
results=("applied", "rejected"),
310+
reasons=(None, "operation_identity_mismatch"),
311+
authority_revision=1,
312+
receipts=("op-A",),
313+
claimed_by=(("todo-1", "agent-a"), ("todo-2", None)),
314+
todo_revisions=(("todo-1", 8), ("todo-2", 7)),
315+
flags=(("state_unchanged_after_rejection", True),),
316+
),
317+
"stale_generation_does_not_duplicate": Expectation(
318+
results=("applied", "applied"),
319+
reasons=(None, None),
320+
authority_revision=2,
321+
receipts=("op-agent-a-todo-1", "op-agent-b-todo-2"),
322+
claimed_by=(("todo-1", "agent-a"), ("todo-2", "agent-b")),
323+
todo_revisions=(("todo-1", 8), ("todo-2", 8)),
324+
flags=(
325+
("stale_expectation_offered", True),
326+
("stale_cas_refused", True),
327+
("stale_cas_wrote_nothing", True),
328+
),
329+
),
330+
}
331+
332+
333+
@pytest.fixture
334+
def handles(tmp_path) -> Callable[[], Any]:
335+
"""Return a factory for fresh provider handles onto one shared store."""
336+
337+
directory = tmp_path / "coordination"
338+
return lambda: FileCoordinationProvider(directory, GOAL_ID)
339+
340+
341+
@pytest.mark.parametrize("name", sorted(SCENARIOS))
342+
def test_claim_work_contract_holds_on_the_shipped_provider(name, handles):
343+
expected = EXPECTED[name]
344+
observation = SCENARIOS[name](handles)
345+
346+
assert tuple(item["result"] for item in observation.outcomes) == expected.results
347+
assert tuple(item.get("reason") for item in observation.outcomes) == expected.reasons
348+
assert observation.head["authority_revision"] == expected.authority_revision
349+
assert tuple(sorted(observation.head["receipt_index"])) == expected.receipts
350+
351+
todos = observation.head["coordination"]["todos"]
352+
assert tuple(
353+
(todo_id, todos[todo_id]["claimed_by"]) for todo_id, _ in expected.claimed_by
354+
) == expected.claimed_by
355+
assert tuple(
356+
(todo_id, todos[todo_id]["todo_revision"])
357+
for todo_id, _ in expected.todo_revisions
358+
) == expected.todo_revisions
359+
360+
for flag, value in expected.flags:
361+
assert observation.flags.get(flag) is value, flag

0 commit comments

Comments
 (0)