diff --git a/loopx/cli_commands/todo_event.py b/loopx/cli_commands/todo_event.py index 102975bc3b..b7ace0bbb2 100644 --- a/loopx/cli_commands/todo_event.py +++ b/loopx/cli_commands/todo_event.py @@ -4,6 +4,9 @@ from collections.abc import Callable from pathlib import Path +from ..control_plane.coordination.local_authority import ( + LocalCoordinationAuthorityRejection, +) from ..control_plane.todos.external_wait_contract import TodoExternalWaitAuthoringError from ..control_plane.todos.handoff_mode import HandoffModeError from ..control_plane.todos.contract import decision_scope_metadata_value @@ -43,12 +46,25 @@ def todo_error_payload(args: argparse.Namespace, exc: Exception) -> dict[str, ob if isinstance(exc, (TaskLeaseError, HandoffModeError)): payload["error_code"] = exc.code payload.update(exc.payload) + elif isinstance(exc, LocalCoordinationAuthorityRejection): + payload["error_code"] = exc.code + payload["code"] = exc.code + for key, value in exc.payload.items(): + if key not in { + "schema_version", + "status", + "failure_kind", + "reason_code", + "reason", + }: + payload[key] = value elif isinstance(exc, TodoExternalWaitAuthoringError): payload["error_code"] = exc.code if exc.authoring_contract is not None: payload["authoring_contract"] = exc.authoring_contract return payload + def append_todo_rollout_event( payload: dict[str, object], *, diff --git a/loopx/control_plane/coordination/local_authority.py b/loopx/control_plane/coordination/local_authority.py index ccc7088cc4..c6a3a1701e 100644 --- a/loopx/control_plane/coordination/local_authority.py +++ b/loopx/control_plane/coordination/local_authority.py @@ -43,6 +43,25 @@ def __init__(self, message: str, *, code: str, payload: Mapping[str, Any]) -> No self.payload = dict(payload) +class LocalCoordinationAuthorityRejection( + LocalCoordinationAuthorityUnavailable, ValueError +): + """The TypeScript coordination owner definitively rejected a claim. + + The legacy Python kernel raised ``ValueError`` for every claim rejection + (todo_not_open, claim_owner_mismatch, unregistered actor, ...). After + promotion those rejections surface as ``status="failed"`` results from the + TypeScript transaction owner; re-raising them through this class keeps the + legacy ``except ValueError`` contract intact for Python API callers while + remaining catchable as an authority outage. Infrastructure and protocol + failures keep raising :class:`LocalCoordinationAuthorityUnavailable`, which + is not a ``ValueError``. + """ + + def __init__(self, message: str, *, code: str, payload: Mapping[str, Any]) -> None: + super().__init__(message, code=code, payload=payload) + + def local_authority_is_promoted(*, runtime_root: Path, goal_id: str) -> bool: fence_path = legacy_coordination_writer_fence_path( runtime_root=runtime_root, @@ -94,7 +113,11 @@ def claim_canonical_todo_if_promoted( "registered_agents": registered_agent_ids_from_registry( registry_path, goal_id ), - "operation_id": operation_id if operation_id is not None else f"todo-claim:{goal_id}:{todo_id}:{uuid4().hex}", + "operation_id": ( + operation_id + if operation_id is not None + else f"todo-claim:{goal_id}:{todo_id}:{uuid4().hex}" + ), "lease_request": ( { "idempotency_key": task_lease_idempotency_key, @@ -116,6 +139,19 @@ def claim_canonical_todo_if_promoted( ) payload = dict(result) accepted = {"applied", "recovered", "replayed", "no_change", "planned"} + if ( + payload.get("status") == "failed" + and payload.get("failure_kind") == "decision_rejection" + ): + # The TypeScript owner classifies this failure as a definitive claim + # decision. The legacy kernel raised ValueError for the same + # rejections, so keep that caller-observable contract; protocol and + # storage-integrity failures stay infrastructure outages. + raise LocalCoordinationAuthorityRejection( + str(payload.get("reason") or "canonical Todo claim was rejected"), + code=str(payload.get("reason_code") or "claim_rejected"), + payload=payload, + ) if ( payload.get("status") not in accepted or payload.get("source_authority") != "file_v0" @@ -188,7 +224,9 @@ def read_canonical_todos_if_promoted( ): raise LocalCoordinationAuthorityUnavailable( str(payload.get("reason") or "canonical Todo authority is unavailable"), - code=str(payload.get("reason_code") or "local_authority_todo_list_unavailable"), + code=str( + payload.get("reason_code") or "local_authority_todo_list_unavailable" + ), payload=payload, ) payload["todos"] = [dict(item) for item in todos] @@ -206,7 +244,8 @@ def canonical_todo_summary_fields( from ..todos.todo_summary import compact_todo_group, count_advancement_todos native_archived = { - item["todo_id"] for item in todos + item["todo_id"] + for item in todos if item.get("schema_version") == TODO_DOMAIN_ITEM_SCHEMA_VERSION and item.get("archive_state") == "archive" } @@ -217,12 +256,14 @@ def canonical_todo_summary_fields( **item, "schema_version": TODO_ITEM_SCHEMA_VERSION, "source_section": ( - "Completed Work Archive" if item["archive_state"] == "archive" + "Completed Work Archive" + if item["archive_state"] == "archive" else TODO_SECTION_HEADINGS[item["role"]] ), "index": index, } - if item.get("schema_version") == TODO_DOMAIN_ITEM_SCHEMA_VERSION else item + if item.get("schema_version") == TODO_DOMAIN_ITEM_SCHEMA_VERSION + else item for index, item in enumerate(todos, 1) ] fields: dict[str, Any] = {} @@ -243,10 +284,14 @@ def canonical_todo_summary_fields( ) if summary: if role == "agent": - archived_done = count_advancement_todos([ - item for item in todos - if item.get("todo_id") in native_archived and item.get("done") is True - ]) + archived_done = count_advancement_todos( + [ + item + for item in todos + if item.get("todo_id") in native_archived + and item.get("done") is True + ] + ) if archived_done: summary["archived_advancement_done_count"] = archived_done summary["advancement_done_count"] = ( diff --git a/loopx/control_plane/coordination/todo_agents.ts b/loopx/control_plane/coordination/todo_agents.ts index bac6bd61de..e84de3a474 100644 --- a/loopx/control_plane/coordination/todo_agents.ts +++ b/loopx/control_plane/coordination/todo_agents.ts @@ -1,10 +1,31 @@ import {AuthorityStoreProtocolError} from "./authority_store_codec.ts"; +// Python str.split() and isspace() recognize exactly 29 Unicode whitespace +// code points: ASCII \t\n\v\f\r and space, ASCII information separators +// U+001C..U+001F, the C1 control NEL (U+0085), and Unicode whitespace blocks +// (NBSP U+00A0, Ogham space mark U+1680, en/em/thin spaces U+2000..U+200A, +// line/paragraph separators U+2028/U+2029, mathematical/ideographic spaces +// U+202F/U+205F/U+3000). Notably, ECMAScript \s omits U+001C..U+001F and +// U+0085 while including BOM (U+FEFF), which Python rejects as whitespace. +// Explicitly match Python's exact 29-code-point whitespace set. +const PYTHON_WHITESPACE_CLASS = + "[\\t\\n\\v\\f\\r \\u001c-\\u001f\\u0085\\u00a0\\u1680\\u2000-\\u200a\\u2028\\u2029\\u202f\\u205f\\u3000]"; +const PYTHON_LEADING_TRAILING_WHITESPACE = new RegExp( + `^${PYTHON_WHITESPACE_CLASS}+|${PYTHON_WHITESPACE_CLASS}+$`, + "gu", +); +const PYTHON_WHITESPACE_RUN = new RegExp(`${PYTHON_WHITESPACE_CLASS}+`, "gu"); + export function normalizeTodoAgent(value: unknown, label: string): string { if (typeof value !== "string") { throw new AuthorityStoreProtocolError(`${label} must be a public-safe agent id`); } - const candidate = value.trim().toLowerCase().replaceAll(" ", "-"); + // Trim and collapse every Python-equivalent whitespace run to one "-" so ids + // typed with any Python whitespace (including U+0085 NEL, U+001C..U+001F, + // tabs, and NBSP) fold exactly like the Python kernel's compact_todo_text path + // (loopx/control_plane/todos/contract.py normalize_todo_claimed_by). + const stripped = value.replace(PYTHON_LEADING_TRAILING_WHITESPACE, ""); + const candidate = stripped.toLowerCase().replace(PYTHON_WHITESPACE_RUN, "-"); if (!/^[a-z][a-z0-9_.:@-]{0,79}$/u.test(candidate)) { throw new AuthorityStoreProtocolError(`${label} must be a public-safe agent id`); } diff --git a/loopx/control_plane/coordination/todo_claim.ts b/loopx/control_plane/coordination/todo_claim.ts index 7ce6e6573c..e412a8de48 100644 --- a/loopx/control_plane/coordination/todo_claim.ts +++ b/loopx/control_plane/coordination/todo_claim.ts @@ -97,11 +97,21 @@ function normalizeExcludedAgents(value: unknown): string[] { return normalized; } -function failure(code: string, reason: string, detail: JsonObject = {}): CoordinationTodoClaimResult { +type CoordinationTodoClaimFailureKind = + | "decision_rejection" + | "protocol_failure"; + +function failure( + code: string, + reason: string, + detail: JsonObject = {}, + kind: CoordinationTodoClaimFailureKind = "protocol_failure", +): CoordinationTodoClaimResult { return { ...detail, schema_version: COORDINATION_TODO_CLAIM_RESULT_SCHEMA, status: "failed", + failure_kind: kind, reason_code: code, reason, }; @@ -422,9 +432,12 @@ export async function executeCoordinationTodoClaim( } const todo = projection.todos.get(input.todo_id); if (todo === undefined) { - return failure("todo_not_found", "Todo is missing from the canonical provider head", { - todo_id: input.todo_id, - }); + return failure( + "todo_not_found", + "Todo is missing from the canonical provider head", + { todo_id: input.todo_id }, + "decision_rejection", + ); } let authority: ReturnType; @@ -441,6 +454,7 @@ export async function executeCoordinationTodoClaim( typeof authority.reason_code === "string" ? authority.reason_code : "invalid_coordination_todo_claim", typeof authority.reason === "string" ? authority.reason : "Todo claim was rejected", authority, + "decision_rejection", ); } @@ -455,6 +469,7 @@ export async function executeCoordinationTodoClaim( "claim_lease_requires_hard_lease", "atomic Todo claim and lease acquire requires handoff_mode=hard_lease", { todo_id: input.todo_id, handoff_mode: handoffMode }, + "decision_rejection", ); } @@ -557,14 +572,34 @@ export async function executeCoordinationTodoClaim( actor_agent_id: authority.owner, lease_decision: decision, }, + "decision_rejection", ); } } else if (handoffMode === "hard_lease" && !activeLeaseForOwner(currentLease, authority.owner, input.now)) { + const existingVersion = currentLease !== undefined + ? (leaseInteger(currentLease, "version") ?? 0) + : null; + const expectedVersionGuidance = existingVersion !== null + ? `; specify --task-lease-expected-version ${existingVersion} to match the existing canonical lease version` + : "; provide --task-lease-expected-version if a canonical lease already exists"; return failure( "handoff_mode_requires_lease", - "hard_lease Todo claim requires an active canonical lease held by the claiming agent", - { todo_id: input.todo_id, actor_agent_id: authority.owner }, + `hard_lease Todo claim requires an active canonical lease held by the claiming agent; ` + + `retry with \`loopx todo claim --task-lease-idempotency-key \`${expectedVersionGuidance}`, + { + todo_id: input.todo_id, + actor_agent_id: authority.owner, + handoff_mode: handoffMode, + recovery: { + command: "loopx todo claim", + requires_flags: ["--task-lease-idempotency-key"], + optional_flags: ["--task-lease-expected-version"], + expected_version: existingVersion, + expected_version_required: existingVersion !== null, + }, + }, + "decision_rejection", ); } } catch (error) { diff --git a/tests/control_plane/test_local_coordination_authority.py b/tests/control_plane/test_local_coordination_authority.py index f4b52c6baf..7624d0b2c7 100644 --- a/tests/control_plane/test_local_coordination_authority.py +++ b/tests/control_plane/test_local_coordination_authority.py @@ -10,6 +10,7 @@ import pytest from loopx.control_plane.coordination.local_authority import ( + LocalCoordinationAuthorityRejection, LocalCoordinationAuthorityUnavailable, claim_canonical_todo_if_promoted, read_canonical_todos_if_promoted, @@ -133,10 +134,13 @@ def test_absent_fence_preserves_legacy_path_without_starting_typescript( "loopx.control_plane.coordination.local_authority.effect_runtime_result", lambda *_args, **_kwargs: pytest.fail("pre-cutover read must stay legacy"), ) - assert read_canonical_todos_if_promoted( - runtime_root=tmp_path, - goal_id="goal-a", - ) is None + assert ( + read_canonical_todos_if_promoted( + runtime_root=tmp_path, + goal_id="goal-a", + ) + is None + ) def test_engaged_fence_reads_typescript_provider_result( @@ -177,9 +181,7 @@ def test_promoted_claim_adapter_invokes_typescript_without_markdown_fallback( "goals": [ { "id": "goal-a", - "coordination": { - "registered_agents": ["agent-a", "agent-b"] - }, + "coordination": {"registered_agents": ["agent-a", "agent-b"]}, } ], } @@ -221,9 +223,7 @@ def _claim(method: str, params: dict[str, object]) -> dict[str, object]: assert result is not None and result["changed"] is True assert calls[0][0] == "coordination.local_authority.todo_claim" assert calls[0][1]["registered_agents"] == ["agent-a", "agent-b"] - assert str(calls[0][1]["operation_id"]).startswith( - "todo-claim:goal-a:todo_a:" - ) + assert str(calls[0][1]["operation_id"]).startswith("todo-claim:goal-a:todo_a:") assert isinstance(calls[0][1]["observed_at"], str) assert calls[0][1]["lease_request"] == { "idempotency_key": "turn:claim-and-acquire", @@ -237,14 +237,21 @@ def test_promoted_add_invokes_native_create_without_markdown_state( tmp_path: Path, ) -> None: registry = tmp_path / "registry.json" - registry.write_text(json.dumps({ - "schema_version": 1, - "common_runtime_root": str(tmp_path / "runtime"), - "goals": [{ - "id": "goal-a", - "coordination": {"registered_agents": ["agent-a", "agent-b"]}, - }], - }), encoding="utf-8") + registry.write_text( + json.dumps( + { + "schema_version": 1, + "common_runtime_root": str(tmp_path / "runtime"), + "goals": [ + { + "id": "goal-a", + "coordination": {"registered_agents": ["agent-a", "agent-b"]}, + } + ], + } + ), + encoding="utf-8", + ) calls: list[tuple[str, dict[str, object]]] = [] monkeypatch.setattr( "loopx.control_plane.todos.provider_create.read_canonical_todos_if_promoted", @@ -254,7 +261,8 @@ def test_promoted_add_invokes_native_create_without_markdown_state( def _create(method: str, params: dict[str, object]) -> dict[str, object]: calls.append((method, params)) return { - "status": "applied", "changed": True, + "status": "applied", + "changed": True, "source_authority": "file_v0", "decision_read_from_provider": True, "legacy_fallback_used": False, @@ -264,9 +272,14 @@ def _create(method: str, params: dict[str, object]) -> dict[str, object]: "loopx.control_plane.todos.provider_create.effect_runtime_result", _create ) result = add_goal_todo( - registry_path=registry, goal_id="goal-a", role="agent", - text="Create natively", claimed_by="agent-a", agent_id="agent-a", - task_class="advancement_task", action_kind="implement", + registry_path=registry, + goal_id="goal-a", + role="agent", + text="Create natively", + claimed_by="agent-a", + agent_id="agent-a", + task_class="advancement_task", + action_kind="implement", validation_command_json='["python", "-c", "pass"]', ) @@ -274,9 +287,7 @@ def _create(method: str, params: dict[str, object]) -> dict[str, object]: assert calls[0][0] == "coordination.local_authority.todo_create" assert calls[0][1]["todo"]["schema_version"] == "todo_domain_record_v0" assert calls[0][1]["todo"]["claimed_by"] == "agent-a" - assert calls[0][1]["todo"]["validation_command_argv"] == [ - "python", "-c", "pass" - ] + assert calls[0][1]["todo"]["validation_command_argv"] == ["python", "-c", "pass"] assert calls[0][1]["registered_agents"] == ["agent-a", "agent-b"] @@ -286,35 +297,55 @@ def test_promoted_add_delegates_semantic_duplicate_to_typescript( ) -> None: monkeypatch.setattr( "loopx.control_plane.todos.provider_create.read_canonical_todos_if_promoted", - lambda **_kwargs: {"todos": [{ - "todo_id": "todo_existing", "role": "agent", "status": "open", - "archive_state": "active", "text": "Already native", - }]}, + lambda **_kwargs: { + "todos": [ + { + "todo_id": "todo_existing", + "role": "agent", + "status": "open", + "archive_state": "active", + "text": "Already native", + } + ] + }, ) calls: list[tuple[str, dict[str, object]]] = [] def _create(method: str, params: dict[str, object]) -> dict[str, object]: calls.append((method, params)) return { - "status": "no_change", "changed": False, - "todo_id": "todo_existing", "source_authority": "file_v0", - "decision_read_from_provider": True, "legacy_fallback_used": False, + "status": "no_change", + "changed": False, + "todo_id": "todo_existing", + "source_authority": "file_v0", + "decision_read_from_provider": True, + "legacy_fallback_used": False, } monkeypatch.setattr( - "loopx.control_plane.todos.provider_create.effect_runtime_result", _create, + "loopx.control_plane.todos.provider_create.effect_runtime_result", + _create, ) registry = tmp_path / "registry.json" - registry.write_text(json.dumps({ - "schema_version": 1, - "goals": [{"id": "goal-a", "coordination": { - "registered_agents": ["agent-a"] - }}], - }), encoding="utf-8") + registry.write_text( + json.dumps( + { + "schema_version": 1, + "goals": [ + {"id": "goal-a", "coordination": {"registered_agents": ["agent-a"]}} + ], + } + ), + encoding="utf-8", + ) result = add_goal_todo( - registry_path=registry, goal_id="goal-a", role="agent", - text="Already native", claimed_by="agent-a", agent_id="agent-a", + registry_path=registry, + goal_id="goal-a", + role="agent", + text="Already native", + claimed_by="agent-a", + agent_id="agent-a", ) assert result["already_exists"] is True @@ -433,6 +464,246 @@ def test_engaged_fence_never_falls_back_when_provider_is_missing( assert exc_info.value.code == "local_authority_todo_list_unavailable" +def _claim_registry(tmp_path: Path) -> Path: + registry = tmp_path / "registry.json" + registry.write_text( + json.dumps( + { + "schema_version": 1, + "goals": [ + { + "id": "goal-a", + "coordination": {"registered_agents": ["agent-a"]}, + } + ], + } + ), + encoding="utf-8", + ) + return registry + + +def _seed_promoted_store( + runtime_root: Path, *, handoff_mode: str | None = None +) -> None: + """Promote one open agent Todo through the real TypeScript runtime.""" + + projection = build_todo_runtime_shadow_projection( + goal_id="goal-a", + todos=[ + { + "schema_version": "todo_item_v0", + "index": 1, + "done": False, + "text": "Claim through the promoted provider head", + "todo_id": "todo_a", + "role": "agent", + "status": "open", + "archive_state": "active", + "source_section": TODO_SECTION_HEADINGS["agent"], + } + ], + ) + if handoff_mode is not None: + projection["handoff_mode"] = str(handoff_mode) + canonical_bytes = json.dumps( + projection, + ensure_ascii=False, + sort_keys=True, + separators=(",", ":"), + ).encode("utf-8") + projection_sha256 = hashlib.sha256(canonical_bytes).hexdigest() + bootstrap = effect_runtime_result( + "coordination.runtime_shadow.bootstrap", + { + "schema_version": "loopx_coordination_runtime_shadow_bootstrap_v0", + "runtime_root": str(runtime_root), + "goal_id": "goal-a", + "operation_id": "bootstrap:goal-a:f34", + "source_version": "state:f34:0", + "projection": projection, + }, + ) + assert bootstrap["status"] == "applied" + mirrored = effect_runtime_result( + "coordination.runtime_shadow.commit", + { + "schema_version": "loopx_coordination_runtime_shadow_commit_v0", + "runtime_root": str(runtime_root), + "goal_id": "goal-a", + "operation_id": "todo:goal-a:f34:qualify", + "event_kind": "todo_update", + "source_version": "state:f34:1", + "projection": projection, + }, + ) + assert mirrored["status"] == "applied" + provider_revision = str(mirrored["provider_revision"]) + fence = { + "schema_version": "loopx_legacy_coordination_writer_fence_v0", + "state": "engaged", + "goal_id": "goal-a", + "fence_id": "legacy-writer-fence:goal-a:f34", + "source_version": "state:f34:1", + "source_projection_sha256": projection_sha256, + "expected_shadow_provider_revision": provider_revision, + } + engaged = effect_runtime_result( + "coordination.local_authority.legacy_writer_fence.engage", + { + "schema_version": "loopx_legacy_coordination_writer_fence_engage_request_v0", + "runtime_root": str(runtime_root), + "goal_id": "goal-a", + "fence": fence, + }, + ) + assert engaged["status"] == "applied" + promoted = effect_runtime_result( + "coordination.local_authority.promote", + { + "schema_version": "loopx_local_coordination_promotion_request_v0", + "runtime_root": str(runtime_root), + "goal_id": "goal-a", + "operation_id": "promote:goal-a:f34", + "expected_shadow_provider_revision": provider_revision, + "expected_shadow_projection_sha256": projection_sha256, + "minimum_operations": 1, + "required_event_kinds": ["todo_update"], + "writer_fence": fence, + }, + ) + assert promoted["status"] == "applied" + + +def test_promoted_claim_rejection_preserves_legacy_valueerror_contract( + monkeypatch: pytest.MonkeyPatch, + tmp_path: Path, +) -> None: + """Promoted claim rejections must stay catchable via ``except ValueError``. + + The legacy kernel raised ValueError for decision rejections such as + todo_not_open; external Python API callers rely on that contract. + """ + _engage_fence(tmp_path) + monkeypatch.setattr( + "loopx.control_plane.coordination.local_authority.effect_runtime_result", + lambda method, params: { + "status": "failed", + "failure_kind": "decision_rejection", + "reason_code": "todo_not_open", + "reason": "todo claim requires status=open", + "source_authority": "file_v0", + "decision_read_from_provider": True, + "legacy_fallback_used": False, + }, + ) + with pytest.raises(ValueError) as exc_info: + claim_canonical_todo_if_promoted( + registry_path=_claim_registry(tmp_path), + runtime_root=tmp_path, + goal_id="goal-a", + todo_id="todo_a", + role="agent", + claimed_by="agent-a", + actor_agent_id="agent-a", + dry_run=False, + ) + rejection = exc_info.value + assert isinstance(rejection, LocalCoordinationAuthorityRejection) + # Callers that already migrated to the authority-unavailable handling of + # promoted claims keep working: the rejection is still its subclass. + assert isinstance(rejection, LocalCoordinationAuthorityUnavailable) + assert rejection.code == "todo_not_open" + assert str(rejection) == "todo claim requires status=open" + + +def test_promoted_claim_protocol_failure_stays_infrastructure_outage( + monkeypatch: pytest.MonkeyPatch, + tmp_path: Path, +) -> None: + """Protocol-level request failures must not masquerade as ValueErrors.""" + _engage_fence(tmp_path) + monkeypatch.setattr( + "loopx.control_plane.coordination.local_authority.effect_runtime_result", + lambda method, params: { + "status": "failed", + "reason_code": "invalid_local_coordination_todo_claim_request", + "reason": "registered_agents must be a JSON array", + "source_authority": "file_v0", + "decision_read_from_provider": True, + "legacy_fallback_used": False, + }, + ) + with pytest.raises(LocalCoordinationAuthorityUnavailable) as exc_info: + claim_canonical_todo_if_promoted( + registry_path=_claim_registry(tmp_path), + runtime_root=tmp_path, + goal_id="goal-a", + todo_id="todo_a", + role="agent", + claimed_by="agent-a", + actor_agent_id="agent-a", + dry_run=False, + ) + assert exc_info.value.code == "invalid_local_coordination_todo_claim_request" + assert not isinstance(exc_info.value, ValueError) + assert not isinstance(exc_info.value, LocalCoordinationAuthorityRejection) + + +def test_promoted_claim_folds_agent_id_whitespace_like_legacy(tmp_path: Path) -> None: + """Every Python whitespace character (including U+0085 NEL, U+001C..U+001F, + tabs, NBSP) in claimed_by must fold to '-' identically before and after promotion. + """ + from loopx.control_plane.todos.contract import normalize_todo_claimed_by + + variants = [ + "Agent A", + "Agent\tA", + "Agent\u0085A", + "Agent\u001cA", + "Agent\u001dA", + "Agent\u001eA", + "Agent\u001fA", + "Agent\u00a0A", + "\u0085 Agent \t A \u001c ", + ] + for variant in variants: + assert normalize_todo_claimed_by(variant) == "agent-a" + + _seed_promoted_store(tmp_path) + result = claim_canonical_todo_if_promoted( + registry_path=_claim_registry(tmp_path), + runtime_root=tmp_path, + goal_id="goal-a", + todo_id="todo_a", + role="agent", + claimed_by="Agent\u0085A", + actor_agent_id="\u0085 Agent \t A \u001c ", + dry_run=False, + ) + assert result is not None and result["ok"] is True + assert result["status"] == "applied" + assert result["claimed_by"] == "agent-a" + + +def test_promoted_claim_missing_todo_rejection_is_valueerror(tmp_path: Path) -> None: + """End-to-end: a real TypeScript decision rejection raises ValueError.""" + _seed_promoted_store(tmp_path) + with pytest.raises(ValueError) as exc_info: + claim_canonical_todo_if_promoted( + registry_path=_claim_registry(tmp_path), + runtime_root=tmp_path, + goal_id="goal-a", + todo_id="todo_missing", + role="agent", + claimed_by="agent-a", + actor_agent_id="agent-a", + dry_run=False, + ) + assert isinstance(exc_info.value, LocalCoordinationAuthorityRejection) + assert exc_info.value.code == "todo_not_found" + + def test_todo_list_uses_provider_after_cutover_even_when_markdown_disagrees( monkeypatch: pytest.MonkeyPatch, tmp_path: Path, @@ -457,9 +728,7 @@ def test_todo_list_uses_provider_after_cutover_even_when_markdown_disagrees( "status": "active", "repo": str(project), "state_file": ".codex/goals/goal-a/ACTIVE_GOAL_STATE.md", - "coordination": { - "registered_agents": ["agent-a", "agent-b"] - }, + "coordination": {"registered_agents": ["agent-a", "agent-b"]}, } ], } @@ -550,7 +819,7 @@ def test_promoted_hard_lease_claim_cli_atomically_acquires_ownership( ) state_file.unlink() - command = [ + base_command = [ sys.executable, "-m", "loopx.cli", @@ -570,6 +839,27 @@ def test_promoted_hard_lease_claim_cli_atomically_acquires_ownership( "agent-a", "--claim-operation-id", "atomic-cli-claim", + ] + initial_failure = subprocess.run( + base_command, + capture_output=True, + text=True, + timeout=30, + ) + assert initial_failure.returncode == 1 + failure_payload = json.loads(initial_failure.stdout) + assert failure_payload["ok"] is False + assert failure_payload["error_code"] == "handoff_mode_requires_lease" + assert failure_payload["handoff_mode"] == "hard_lease" + assert "loopx todo claim --task-lease-idempotency-key" in failure_payload["error"] + assert "--task-lease-expected-version" in failure_payload["error"] + recovery = failure_payload.get("recovery") or {} + assert recovery.get("command") == "loopx todo claim" + assert recovery.get("requires_flags") == ["--task-lease-idempotency-key"] + assert "--task-lease-expected-version" in (recovery.get("optional_flags") or []) + + command = [ + *base_command, "--task-lease-idempotency-key", "turn:atomic-cli-claim", "--task-lease-expected-version", @@ -643,9 +933,7 @@ def test_real_shadow_projection_promotes_complete_complex_todo_semantics( "status": "active", "repo": str(project), "state_file": ".codex/goals/goal-a/ACTIVE_GOAL_STATE.md", - "coordination": { - "registered_agents": ["agent-a", "agent-b"] - }, + "coordination": {"registered_agents": ["agent-a", "agent-b"]}, } ], } @@ -755,17 +1043,32 @@ def test_real_shadow_projection_promotes_complete_complex_todo_semantics( # The public compatibility CLI must retain claim-neutral text correction # after promotion; it must not reconstruct or write the Markdown source. correction_command = [ - sys.executable, "-m", "loopx.cli", "--format", "json", - "--registry", str(registry_path), "todo", "update", "--goal-id", "goal-a", - "--todo-id", "todo_claimable", "--agent-id", "agent-b", - "--text", "Corrected before claiming", + sys.executable, + "-m", + "loopx.cli", + "--format", + "json", + "--registry", + str(registry_path), + "todo", + "update", + "--goal-id", + "goal-a", + "--todo-id", + "todo_claimable", + "--agent-id", + "agent-b", + "--text", + "Corrected before claiming", ] - correction = subprocess.run(correction_command, capture_output=True, text=True, - check=True, timeout=30) + correction = subprocess.run( + correction_command, capture_output=True, text=True, check=True, timeout=30 + ) assert json.loads(correction.stdout)["ok"] is True corrected = list_goal_todos(registry_path=registry_path, goal_id="goal-a") - corrected_item = next(item for item in corrected["todos"] - if item["todo_id"] == "todo_claimable") + corrected_item = next( + item for item in corrected["todos"] if item["todo_id"] == "todo_claimable" + ) assert corrected_item["text"] == "Corrected before claiming" assert not corrected_item.get("claimed_by") assert corrected_item["last_actor_agent_id"] == "agent-b" @@ -788,19 +1091,45 @@ def test_real_shadow_projection_promotes_complete_complex_todo_semantics( assert not state_file.exists() claim_command = [ - sys.executable, "-m", "loopx.cli", "--format", "json", - "--registry", str(registry_path), "todo", "claim", "--goal-id", "goal-a", - "--todo-id", "todo_claimable", "--claimed-by", "agent-a", "--agent-id", "agent-a", - "--claim-operation-id", "initial-cli-claim", + sys.executable, + "-m", + "loopx.cli", + "--format", + "json", + "--registry", + str(registry_path), + "todo", + "claim", + "--goal-id", + "goal-a", + "--todo-id", + "todo_claimable", + "--claimed-by", + "agent-a", + "--agent-id", + "agent-a", + "--claim-operation-id", + "initial-cli-claim", ] # Duplicate callers race from separate processes, but one operation must # produce exactly one accepted claim and the same durable receipt. with ThreadPoolExecutor(max_workers=2) as pool: - attempts = [pool.submit(subprocess.run, claim_command, - capture_output=True, text=True, check=True, timeout=30) for _ in range(2)] + attempts = [ + pool.submit( + subprocess.run, + claim_command, + capture_output=True, + text=True, + check=True, + timeout=30, + ) + for _ in range(2) + ] claims = [json.loads(attempt.result().stdout) for attempt in attempts] assert sum(item["status"] == "applied" for item in claims) == 1 - assert all(item["status"] in {"applied", "recovered", "replayed"} for item in claims) + assert all( + item["status"] in {"applied", "recovered", "replayed"} for item in claims + ) assert claims[0]["original_receipt"] == claims[1]["original_receipt"] assert claims[0]["provider_revision"] == claims[1]["provider_revision"] claimed = next(item for item in claims if item["status"] == "applied") @@ -818,44 +1147,90 @@ def test_real_shadow_projection_promotes_complete_complex_todo_semantics( # Separate CLI processes must replay one durable operation, not mint a # fresh receipt for every retry. Preview does not consume that identity. claim_command = [*claim_command[:-1], "retryable-cli-claim"] - preview = json.loads(subprocess.run( - [*claim_command, "--dry-run"], capture_output=True, text=True, check=True, - ).stdout) + preview = json.loads( + subprocess.run( + [*claim_command, "--dry-run"], + capture_output=True, + text=True, + check=True, + ).stdout + ) assert preview["dry_run"] is True - original = json.loads(subprocess.run( - claim_command, capture_output=True, text=True, check=True, - ).stdout) - replay = json.loads(subprocess.run( - claim_command, capture_output=True, text=True, check=True, - ).stdout) + original = json.loads( + subprocess.run( + claim_command, + capture_output=True, + text=True, + check=True, + ).stdout + ) + replay = json.loads( + subprocess.run( + claim_command, + capture_output=True, + text=True, + check=True, + ).stdout + ) assert original["status"] == "no_change" assert replay["status"] == "replayed" assert replay["original_receipt"] == original["original_receipt"] assert replay["provider_revision"] == original["provider_revision"] - changed_intent = ["agent-b" if part == "agent-a" else part for part in claim_command] + changed_intent = [ + "agent-b" if part == "agent-a" else part for part in claim_command + ] rejected = subprocess.run(changed_intent, capture_output=True, text=True) assert rejected.returncode != 0 - assert json.loads(rejected.stdout)["error"] == "operation id already names a different coordination request" + assert ( + json.loads(rejected.stdout)["error"] + == "operation id already names a different coordination request" + ) for invalid_key in ("", " padded-operation "): - invalid = subprocess.run([*claim_command[:-1], invalid_key], capture_output=True, text=True) + invalid = subprocess.run( + [*claim_command[:-1], invalid_key], capture_output=True, text=True + ) assert invalid.returncode != 0 assert not state_file.exists() create_command = [ - sys.executable, "-m", "loopx.cli", "--format", "json", - "--registry", str(registry_path), "todo", "add", "--goal-id", "goal-a", - "--role", "agent", "--text", "Create directly against promoted provider", - "--claimed-by", "agent-a", - "--task-class", "advancement_task", "--action-kind", "implement", + sys.executable, + "-m", + "loopx.cli", + "--format", + "json", + "--registry", + str(registry_path), + "todo", + "add", + "--goal-id", + "goal-a", + "--role", + "agent", + "--text", + "Create directly against promoted provider", + "--claimed-by", + "agent-a", + "--task-class", + "advancement_task", + "--action-kind", + "implement", ] create_preview = subprocess.run( - [*create_command, "--dry-run"], capture_output=True, text=True, check=True, + [*create_command, "--dry-run"], + capture_output=True, + text=True, + check=True, ) assert json.loads(create_preview.stdout)["status"] == "planned" assert not state_file.exists() - created = json.loads(subprocess.run( - create_command, capture_output=True, text=True, check=True, - ).stdout) + created = json.loads( + subprocess.run( + create_command, + capture_output=True, + text=True, + check=True, + ).stdout + ) assert created["ok"] is True assert created["source_authority"] == "file_v0" assert created["legacy_fallback_used"] is False @@ -873,13 +1248,35 @@ def test_real_shadow_projection_promotes_complete_complex_todo_semantics( # Real CLI, no Markdown file: provider data feeds an in-memory editor and # only requested fields return through TS CAS. Complex sibling fields do # not round-trip through the lossy Markdown representation. - command = [sys.executable, "-m", "loopx.cli", "--format", "json", - "--registry", str(registry_path), "todo", "update", "--goal-id", "goal-a", - "--todo-id", "todo_claimable", "--agent-id", "agent-a", - "--text", "Edit provider-owned work", "--note", "compatibility edit"] - preview = subprocess.run([*command, "--dry-run"], capture_output=True, text=True, check=True) + command = [ + sys.executable, + "-m", + "loopx.cli", + "--format", + "json", + "--registry", + str(registry_path), + "todo", + "update", + "--goal-id", + "goal-a", + "--todo-id", + "todo_claimable", + "--agent-id", + "agent-a", + "--text", + "Edit provider-owned work", + "--note", + "compatibility edit", + ] + preview = subprocess.run( + [*command, "--dry-run"], capture_output=True, text=True, check=True + ) assert json.loads(preview.stdout)["status"] == "planned" - assert list_goal_todos(registry_path=registry_path, goal_id="goal-a")["todos"] == after_create["todos"] + assert ( + list_goal_todos(registry_path=registry_path, goal_id="goal-a")["todos"] + == after_create["todos"] + ) edited = subprocess.run(command, capture_output=True, text=True, check=True) edit_result = json.loads(edited.stdout) assert edit_result["status"] == "applied" @@ -888,10 +1285,187 @@ def test_real_shadow_projection_promotes_complete_complex_todo_semantics( after_edit = list_goal_todos(registry_path=registry_path, goal_id="goal-a") edited_by_id = {item["todo_id"]: item for item in after_edit["todos"]} assert edited_by_id["todo_claimable"] == { - **claimed_item, "text": "Edit provider-owned work", "note": "compatibility edit", + **claimed_item, + "text": "Edit provider-owned work", + "note": "compatibility edit", "last_actor_agent_id": "agent-a", "updated_at": edited_by_id["todo_claimable"]["updated_at"], } assert edited_by_id["todo_claimable"]["updated_at"] != claimed_item["updated_at"] assert edited_by_id["todo_complex"] == by_id["todo_complex"] assert edited_by_id["todo_successor"] == by_id["todo_successor"] + + +@pytest.mark.parametrize( + ("reason_code", "reason"), + [ + ( + "coordination_operation_identity_mismatch", + "operation id names a different coordination request", + ), + ( + "invalid_coordination_projection", + "canonical projection could not be decoded", + ), + ( + "invalid_coordination_todo_claim_receipt", + "stored claim receipt could not be verified", + ), + ], +) +def test_promoted_claim_protocol_failures_stay_unavailable( + monkeypatch: pytest.MonkeyPatch, + tmp_path: Path, + reason_code: str, + reason: str, +) -> None: + """Protocol and storage-integrity failures are outages, not decisions.""" + + _engage_fence(tmp_path) + + def _protocol_failure(*_args: object, **_kwargs: object) -> dict[str, object]: + return { + "schema_version": "loopx_coordination_todo_claim_result_v0", + "status": "failed", + "failure_kind": "protocol_failure", + "reason_code": reason_code, + "reason": reason, + } + + monkeypatch.setattr( + "loopx.control_plane.coordination.local_authority.effect_runtime_result", + _protocol_failure, + ) + with pytest.raises(LocalCoordinationAuthorityUnavailable) as exc_info: + claim_canonical_todo_if_promoted( + registry_path=_claim_registry(tmp_path), + runtime_root=tmp_path, + goal_id="goal-a", + todo_id="todo_x1", + role="agent", + claimed_by="agent-a", + actor_agent_id="agent-a", + dry_run=False, + ) + assert not isinstance(exc_info.value, ValueError) + assert exc_info.value.code == reason_code + + +def test_protocol_failure_kind_absent_means_unavailable( + monkeypatch: pytest.MonkeyPatch, + tmp_path: Path, +) -> None: + """A legacy failure result without the kind field is an outage, not a decision.""" + + _engage_fence(tmp_path) + + def _legacy_failure(*_args: object, **_kwargs: object) -> dict[str, object]: + return { + "schema_version": "loopx_coordination_todo_claim_result_v0", + "status": "failed", + "reason_code": "todo_not_open", + "reason": "todo claim requires status=open", + } + + monkeypatch.setattr( + "loopx.control_plane.coordination.local_authority.effect_runtime_result", + _legacy_failure, + ) + with pytest.raises(LocalCoordinationAuthorityUnavailable) as exc_info: + claim_canonical_todo_if_promoted( + registry_path=_claim_registry(tmp_path), + runtime_root=tmp_path, + goal_id="goal-a", + todo_id="todo_a", + role="agent", + claimed_by="agent-a", + actor_agent_id="agent-a", + dry_run=False, + ) + assert not isinstance(exc_info.value, ValueError) + + +def test_hard_lease_eligibility_rejection_is_valueerror( + tmp_path: Path, +) -> None: + """A real hard-lease Todo without a lease rejects with a usable repair path. + + The unpromoted path raises TaskLeaseError (a ValueError) when a hard-lease + Todo has no matching active lease; the promoted path must keep that caller + contract and leave the canonical state untouched. The canonical recovery + operation is to supply a task-lease idempotency key to acquire the lease + atomically with the claim. + """ + + _seed_promoted_store(tmp_path, handoff_mode="hard_lease") + registry_path = _claim_registry(tmp_path) + store_head = json.loads( + ( + tmp_path / "authority" / "file-v0" / "authority-store-bf21e67b01a351a1.json" + ).read_text(encoding="utf-8") + ) + head_before = json.dumps(store_head.get("head", {}), sort_keys=True) + with pytest.raises(ValueError) as exc_info: + claim_canonical_todo_if_promoted( + registry_path=registry_path, + runtime_root=tmp_path, + goal_id="goal-a", + todo_id="todo_a", + role="agent", + claimed_by="agent-a", + actor_agent_id="agent-a", + dry_run=False, + ) + assert isinstance(exc_info.value, LocalCoordinationAuthorityRejection) + assert exc_info.value.code == "handoff_mode_requires_lease" + assert ( + "hard_lease Todo claim requires an active canonical lease held by the claiming agent" + in str(exc_info.value) + ) + assert "loopx todo claim --task-lease-idempotency-key" in str(exc_info.value) + assert "--task-lease-expected-version" in str(exc_info.value) + assert ( + exc_info.value.payload.get("recovery", {}).get("requires_flags") + == ["--task-lease-idempotency-key"] + ) + store_after = json.loads( + ( + tmp_path / "authority" / "file-v0" / "authority-store-bf21e67b01a351a1.json" + ).read_text(encoding="utf-8") + ) + assert json.dumps(store_after.get("head", {}), sort_keys=True) == head_before + + # Canonical recovery path: supply task_lease_idempotency_key to acquire + # the lease atomically during claim. + recovered = claim_canonical_todo_if_promoted( + registry_path=registry_path, + runtime_root=tmp_path, + goal_id="goal-a", + todo_id="todo_a", + role="agent", + claimed_by="agent-a", + actor_agent_id="agent-a", + dry_run=False, + task_lease_idempotency_key="turn:claim-and-acquire", + task_lease_expected_version=0, + ) + assert recovered is not None and recovered["status"] == "applied" + assert recovered["changed"] is True + assert recovered["lease"]["owner"] == "agent-a" + assert recovered["lease"]["status"] == "active" + assert recovered["lease"]["idempotency_key"] == "turn:claim-and-acquire" + + # With the active lease now persisted, subsequent claims succeed without + # requiring a new lease request. + subsequent = claim_canonical_todo_if_promoted( + registry_path=registry_path, + runtime_root=tmp_path, + goal_id="goal-a", + todo_id="todo_a", + role="agent", + claimed_by="agent-a", + actor_agent_id="agent-a", + dry_run=False, + ) + assert subsequent is not None + assert subsequent["status"] in {"applied", "no_change", "replayed"} diff --git a/tests/control_plane_ts/local_authority_runtime.test.ts b/tests/control_plane_ts/local_authority_runtime.test.ts index f45291e3f4..3afb3e5b4f 100644 --- a/tests/control_plane_ts/local_authority_runtime.test.ts +++ b/tests/control_plane_ts/local_authority_runtime.test.ts @@ -7,7 +7,11 @@ import test from "node:test"; import { FileAuthorityStore } from "../../loopx/control_plane/coordination/file_authority_store.ts"; import type { AuthorityStoreCommit } from "../../loopx/control_plane/coordination/authority_store.ts"; -import { canonicalAuthorityBytes } from "../../loopx/control_plane/coordination/authority_store_codec.ts"; +import { + AuthorityStoreProtocolError, + canonicalAuthorityBytes, +} from "../../loopx/control_plane/coordination/authority_store_codec.ts"; +import { normalizeTodoAgent } from "../../loopx/control_plane/coordination/todo_agents.ts"; import { TODO_DOMAIN_ITEM_SCHEMA, TODO_DOMAIN_READ_RECORD_SCHEMA, @@ -512,6 +516,73 @@ test("provider-first Todo claim preserves the complete record and is replay-safe assert.equal(repeated.changed, false); }); +test("agent id normalization folds any whitespace run like the Python kernel", async () => { + // Parity with loopx/control_plane/todos/contract.py normalize_todo_claimed_by: + // compact_todo_text collapses every Python-recognized whitespace run (including + // U+0085 NEL, U+001C..U+001F information separators, tabs, and NBSP) into + // one space before mapping it to "-", so the same claim command keeps working + // before and after promotion. + const whitespaceVariants = [ + "Agent A", + "Agent\tA", + "Agent\u0085A", + "Agent\u001cA", + "Agent\u001dA", + "Agent\u001eA", + "Agent\u001fA", + "Agent\u00a0A", + "\u0085 Agent \t A \u001c ", + ]; + for (const variant of whitespaceVariants) { + assert.equal(normalizeTodoAgent(variant, "claimed_by"), "agent-a"); + } + // BOM (U+FEFF) is not Python whitespace and must not be accepted + assert.throws( + () => normalizeTodoAgent("\ufeffAgent A", "claimed_by"), + AuthorityStoreProtocolError, + ); + + const root = await mkdtemp(join(tmpdir(), "loopx-local-authority-claim-tab-")); + const store = new FileAuthorityStore(join(root, "authority", "file-v0"), "goal-a"); + const seeded = await store.commitAuthority({ + expected_provider_revision: null, + operation_id: "promote:claim-tab-test", + events: [{ schema_version: "promotion_v0" }], + next_projection: withTodoReadModel({ + goal_id: "goal-a", + handoff_mode: "soft_claim", + todos: [todoRecord()], + leases: [], + }), + receipts: [], + }); + assert.equal(seeded.status, "applied"); + + const applied = await claimLocalCoordinationTodo({ + schema_version: LOCAL_COORDINATION_TODO_CLAIM_REQUEST_SCHEMA, + runtime_root: root, + goal_id: "goal-a", + todo_id: "todo_a", + role: "agent", + claimed_by: "Agent\u0085A", + actor_agent_id: "\u0085 Agent \t A \u001c ", + registered_agents: ["agent-a"], + operation_id: "todo-claim:goal-a:todo_a:tab", + observed_at: "2026-09-05T04:30:00Z", + dry_run: false, + }); + assert.equal(applied.status, "applied", JSON.stringify(applied)); + + const read = await readLocalCoordinationTodo({ + schema_version: LOCAL_COORDINATION_TODO_READ_REQUEST_SCHEMA, + runtime_root: root, + goal_id: "goal-a", + todo_id: "todo_a", + }); + assert.equal(read.status, "found"); + assert.equal((read.todo as Record).claimed_by, "agent-a"); +}); + test("one TypeScript decision owns promoted and legacy Todo claims", () => { const input = { goal_id: "goal-a", @@ -855,6 +926,22 @@ test("provider-first Todo claim validates authority and hard-lease ownership", a claimed_by: "agent-a", }); assert.equal(missingLease.reason_code, "handoff_mode_requires_lease"); + assert.match( + String(missingLease.reason), + /loopx todo claim --task-lease-idempotency-key/, + ); + assert.match( + String(missingLease.reason), + /--task-lease-expected-version/, + ); + assert.equal( + (missingLease.recovery as Record | undefined)?.command, + "loopx todo claim", + ); + assert.deepEqual( + (missingLease.recovery as Record | undefined)?.requires_flags, + ["--task-lease-idempotency-key"], + ); const dryRun = await claimLocalCoordinationTodo({ ...base,