"""Client/session-aware runtime ownership and provenance (#948). Covers the reproduced contradiction that motivated the issue: one surface reporting ``client_managed`` while another reported ``manual_launch`` for the same process, remediation hardcoded to one vendor, and a profile-wide duplicate wall that could not tell two healthy clients apart. All client and session identifiers here are synthetic. """ from __future__ import annotations import os import tempfile import unittest from datetime import datetime, timedelta, timezone import mcp_client_reconnect import mcp_namespace_health import mcp_worker_identity as mwi NOW = datetime(2026, 7, 29, 6, 0, 0, tzinfo=timezone.utc) def _registry() -> mwi.WorkerRegistry: """A registry on a throwaway path; never the operator's real one.""" handle, path = tempfile.mkstemp(suffix=".sqlite3") os.close(handle) os.unlink(path) return mwi.WorkerRegistry(path) def _attach( registry: mwi.WorkerRegistry, *, client: str, session: str, generation: str, profile: str = "prgs-reviewer", role: str = "reviewer", pid: int = 4242, now: datetime = NOW, ttl: float = 900.0, ) -> dict: """Register one synthetic worker and return the outcome.""" identity = mwi.generate_worker_identity(client, session, now=now) outcome = registry.register( worker_identity=identity, client_name=client, client_instance_id=f"inst-{session}", session_id=session, generation_id=generation, role=role, profile=profile, pid=pid, heartbeat_ttl_seconds=ttl, now=now, ) outcome["identity"] = identity return outcome class IdentityFormatTests(unittest.TestCase): """AC27-29: collision-resistant `--`.""" def test_identity_matches_required_format(self): identity = mwi.generate_worker_identity("Gemini", "sess-0001", now=NOW) parsed = mwi.parse_worker_identity(identity) self.assertTrue(parsed["valid"], parsed["reasons"]) self.assertEqual(parsed["client_name"], "gemini") self.assertEqual(parsed["minted_at"], "20260729T060000Z") self.assertEqual(len(parsed["digest"]), 12) def test_digest_varies_with_session_and_nonce(self): base = dict(timestamp_ns=1, now=NOW) a = mwi.generate_worker_identity("codex", "sess-A", nonce="n", **base) b = mwi.generate_worker_identity("codex", "sess-B", nonce="n", **base) c = mwi.generate_worker_identity("codex", "sess-A", nonce="m", **base) self.assertNotEqual(a, b, "session must feed the digest") self.assertNotEqual(a, c, "nonce must feed the digest") def test_identity_is_not_role_or_profile(self): """AC26: identity is independent of role and profile.""" args = dict(timestamp_ns=7, nonce="fixed", now=NOW) same = mwi.generate_worker_identity("claude", "sess-1", **args) self.assertEqual(same, mwi.generate_worker_identity("claude", "sess-1", **args)) # Nothing role- or profile-derived appears in the identity. self.assertNotIn("reviewer", same) self.assertNotIn("prgs", same) def test_malformed_identity_rejected(self): self.assertFalse(mwi.parse_worker_identity("prgs-reviewer")["valid"]) self.assertFalse(mwi.parse_worker_identity("")["valid"]) self.assertFalse(mwi.parse_worker_identity(None)["valid"]) class PerClientAttachmentTests(unittest.TestCase): """Every supported client attaches and is reported as itself.""" def _assert_attached_as(self, client: str, expected_name: str): registry = _registry() outcome = _attach( registry, client=client, session=f"sess-{client}", generation="gen-1" ) self.assertTrue(outcome["registered"], outcome["reasons"]) verdict = mwi.assess_provenance( registry=registry, worker_identity=outcome["identity"], env={}, now=NOW ) self.assertEqual(verdict["session_ownership"], mwi.OWNERSHIP_OWNED) self.assertEqual(verdict["provenance"], mwi.PROVENANCE_CLIENT_SESSION) self.assertEqual(verdict["client_name"], expected_name) self.assertTrue(verdict["session_owned"]) self.assertFalse(verdict["fail_closed"]) return verdict def test_codex_attachment(self): self._assert_attached_as("codex", "codex") def test_gemini_attachment(self): self._assert_attached_as("gemini", "gemini") def test_antigravity_attachment(self): self._assert_attached_as("antigravity", "antigravity") def test_claude_attachment(self): self._assert_attached_as("claude", "claude_code") def test_unknown_client_is_named_not_guessed(self): verdict = self._assert_attached_as("some_new_llm", "some_new_llm") self.assertNotEqual(verdict["client_name"], "codex") class SessionLifecycleTests(unittest.TestCase): def test_same_client_new_session_gets_distinct_identity(self): registry = _registry() first = _attach(registry, client="codex", session="sess-1", generation="gen-1") second = _attach(registry, client="codex", session="sess-2", generation="gen-2") self.assertTrue(first["registered"]) self.assertTrue(second["registered"]) self.assertNotEqual(first["identity"], second["identity"]) # Both are live and neither blocks the other. cohort = mwi.classify_cohort(registry.list_workers(), now=NOW) self.assertEqual(cohort["live_worker_count"], 2) self.assertFalse(cohort["blocked"], cohort["reasons"]) def test_different_client_attaches_after_previous_session_ends(self): """AC14: expiry then takeover with a higher fencing epoch.""" registry = _registry() gone = _attach( registry, client="codex", session="sess-old", generation="gen-shared", ttl=60 ) later = NOW + timedelta(hours=1) self.assertFalse( registry.is_live(registry.get(gone["identity"]), now=later)["live"] ) arriving = _attach( registry, client="gemini", session="sess-new", generation="gen-other", now=later, ) claim = registry.claim_generation( worker_identity=arriving["identity"], generation_id="gen-shared", now=later, ) self.assertTrue(claim["claimed"], claim["reasons"]) self.assertIn(gone["identity"], claim["superseded_workers"]) self.assertGreater(claim["fencing_epoch"], gone["fencing_epoch"]) def test_superseded_session_is_fenced_on_resume(self): """AC15/AC16: the prior session cannot heartbeat its way back.""" registry = _registry() old = _attach( registry, client="codex", session="sess-old", generation="gen-shared", ttl=60 ) later = NOW + timedelta(hours=1) new = _attach( registry, client="gemini", session="sess-new", generation="gen-x", now=later ) registry.claim_generation( worker_identity=new["identity"], generation_id="gen-shared", now=later ) resumed = registry.heartbeat( worker_identity=old["identity"], fencing_epoch=old["fencing_epoch"], now=later, ) self.assertFalse(resumed["renewed"]) self.assertFalse(resumed["mutation_performed"]) self.assertEqual(resumed["blocker_kind"], mwi.BLOCKER_FENCED) def test_heartbeat_renews_only_the_owning_lease(self): """AC11: a wrong epoch never renews, and never mutates.""" registry = _registry() worker = _attach(registry, client="codex", session="s", generation="g") good = registry.heartbeat( worker_identity=worker["identity"], fencing_epoch=worker["fencing_epoch"], now=NOW + timedelta(minutes=5), ) self.assertTrue(good["renewed"]) bad = registry.heartbeat( worker_identity=worker["identity"], fencing_epoch=worker["fencing_epoch"] + 99, now=NOW + timedelta(minutes=6), ) self.assertFalse(bad["renewed"]) self.assertFalse(bad["mutation_performed"]) self.assertEqual( registry.get(worker["identity"])["last_heartbeat_at"], good["last_heartbeat_at"], "a refused heartbeat must not advance the record", ) class ConflictAndCollisionTests(unittest.TestCase): def test_two_live_sessions_cannot_claim_one_generation(self): registry = _registry() first = _attach(registry, client="codex", session="s1", generation="gen-shared") second = _attach(registry, client="gemini", session="s2", generation="gen-other") claim = registry.claim_generation( worker_identity=second["identity"], generation_id="gen-shared", now=NOW, ) self.assertFalse(claim["claimed"]) self.assertFalse(claim["mutation_performed"]) self.assertEqual(claim["blocker_kind"], mwi.BLOCKER_CONFLICTING_SESSIONS) self.assertEqual( claim["conflicting_owners"][0]["worker_identity"], first["identity"] ) # The sanctioned recovery must never be "kill the other process". self.assertIn("Do not kill", claim["exact_next_action"]) def test_contested_generation_fails_closed_in_assessment(self): registry = _registry() first = _attach(registry, client="codex", session="s1", generation="gen-shared") _attach(registry, client="gemini", session="s2", generation="gen-shared") verdict = mwi.assess_provenance( registry=registry, worker_identity=first["identity"], env={}, now=NOW ) self.assertEqual(verdict["session_ownership"], mwi.OWNERSHIP_CONTESTED) self.assertTrue(verdict["fail_closed"]) self.assertEqual(verdict["blocker_kind"], mwi.BLOCKER_CONTRADICTORY) self.assertTrue(verdict["conflicting_live_sessions"]) def test_identity_collision_is_refused_without_corrupting_existing(self): """AC31: never replace, adopt, merge with, or corrupt the incumbent.""" registry = _registry() incumbent = _attach(registry, client="codex", session="s1", generation="gen-1") before = registry.get(incumbent["identity"]) collided = registry.register( worker_identity=incumbent["identity"], client_name="gemini", client_instance_id="inst-other", session_id="s2", generation_id="gen-2", pid=9999, now=NOW, ) self.assertFalse(collided["registered"]) self.assertTrue(collided["collision"]) self.assertFalse(collided["mutation_performed"]) self.assertEqual(collided["blocker_kind"], mwi.BLOCKER_IDENTITY_COLLISION) self.assertEqual(collided["collision_kind"], "active_worker") self.assertEqual( registry.get(incumbent["identity"]), before, "incumbent must be untouched" ) def test_after_collision_a_regenerated_identity_registers(self): """AC32/AC35: forced collision, safe regeneration, successful replacement.""" registry = _registry() fixed = dict(timestamp_ns=99, nonce="deterministic", now=NOW) forced = mwi.generate_worker_identity("codex", "sess-collide", **fixed) first = registry.register( worker_identity=forced, client_name="codex", client_instance_id="inst-1", session_id="sess-collide", generation_id="gen-1", now=NOW, ) self.assertTrue(first["registered"]) # A second worker deriving the same inputs collides deterministically. again = mwi.generate_worker_identity("codex", "sess-collide", **fixed) self.assertEqual(again, forced) self.assertTrue( registry.register( worker_identity=again, client_name="codex", client_instance_id="inst-2", session_id="sess-collide", generation_id="gen-2", now=NOW, )["collision"] ) replacement = mwi.generate_worker_identity( "codex", "sess-collide", timestamp_ns=100, nonce="different", now=NOW ) self.assertNotEqual(replacement, forced) self.assertTrue( registry.register( worker_identity=replacement, client_name="codex", client_instance_id="inst-2", session_id="sess-collide", generation_id="gen-2", now=NOW, )["registered"] ) def test_restarted_worker_inherits_nothing(self): """AC33/AC34: a restart mints a new identity and no prior epoch.""" registry = _registry() before = _attach( registry, client="codex", session="sess-before", generation="gen-1", ttl=60 ) later = NOW + timedelta(hours=2) after = _attach( registry, client="codex", session="sess-after", generation="gen-2", now=later ) self.assertNotEqual(before["identity"], after["identity"]) self.assertNotEqual( registry.get(after["identity"])["generation_id"], registry.get(before["identity"])["generation_id"], ) class LivenessTests(unittest.TestCase): def test_stale_session_record_is_not_live(self): registry = _registry() worker = _attach(registry, client="codex", session="s", generation="g", ttl=300) stale = registry.is_live( registry.get(worker["identity"]), now=NOW + timedelta(hours=1) ) self.assertFalse(stale["live"]) self.assertFalse(stale["heartbeat_fresh"]) def test_liveness_is_not_pid_comparison_alone(self): """AC7: a live PID does not resurrect an expired registration.""" registry = _registry() worker = _attach(registry, client="codex", session="s", generation="g", ttl=60) verdict = registry.is_live( registry.get(worker["identity"]), now=NOW + timedelta(hours=1), pid_alive=True, ) self.assertFalse( verdict["live"], "a live PID must not override a dead heartbeat" ) def test_dead_pid_withdraws_liveness_from_a_fresh_heartbeat(self): registry = _registry() worker = _attach(registry, client="codex", session="s", generation="g") verdict = registry.is_live( registry.get(worker["identity"]), now=NOW, pid_alive=False ) self.assertFalse(verdict["live"]) def test_stale_ownership_does_not_permanently_strand_a_daemon(self): registry = _registry() stranded = _attach( registry, client="codex", session="s-old", generation="gen-daemon", ttl=60 ) later = NOW + timedelta(hours=3) rescuer = _attach( registry, client="claude", session="s-new", generation="gen-tmp", now=later ) claim = registry.claim_generation( worker_identity=rescuer["identity"], generation_id="gen-daemon", now=later, ) self.assertTrue(claim["claimed"], claim["reasons"]) self.assertIn(stranded["identity"], claim["superseded_workers"]) class EvidenceTests(unittest.TestCase): def test_env_flag_alone_does_not_prove_session_ownership(self): verdict = mwi.assess_provenance( registry=None, worker_identity=None, env={"GITEA_CLIENT_MANAGED": "1", "GITEA_MCP_SANCTIONED_DAEMON": "1"}, now=NOW, ) self.assertFalse(verdict["session_owned"]) self.assertEqual(verdict["session_ownership"], mwi.OWNERSHIP_UNOWNED) self.assertTrue(verdict["env_flag_only"]) self.assertTrue(verdict["fail_closed"]) self.assertNotIn(mwi.EVIDENCE_ATTACHMENT_RECORD, verdict["evidence"]) self.assertFalse(verdict["env_signal"]["proves_session_ownership"]) def test_env_flag_still_answers_the_launch_question(self): """The #686 wall is preserved: env decides launch, not ownership.""" self.assertTrue( mwi.assess_launch_provenance({"GITEA_CLIENT_MANAGED": "1"})["client_managed"] ) self.assertFalse( mwi.assess_launch_provenance({"GITEA_CLIENT_MANAGED": "0"})["client_managed"] ) self.assertFalse( mwi.assess_launch_provenance({}, stdin_is_tty=True)["client_managed"] ) self.assertTrue( mwi.assess_launch_provenance({"GITEA_MCP_PROFILE": "prgs-author"})[ "client_managed" ] ) def test_missing_evidence_is_unproven_not_manual(self): """A missing proof must not be reported as a hand-launched process.""" verdict = mwi.assess_provenance( registry=None, worker_identity=None, env={}, now=NOW ) self.assertEqual(verdict["provenance"], mwi.PROVENANCE_UNPROVEN) self.assertNotEqual(verdict["provenance"], mwi.PROVENANCE_MANUAL) self.assertTrue(verdict["fail_closed"]) def test_declared_manual_launch_is_reported_as_manual(self): verdict = mwi.assess_provenance( registry=None, worker_identity=None, env={"GITEA_CLIENT_MANAGED": "0"}, now=NOW, ) self.assertEqual(verdict["provenance"], mwi.PROVENANCE_MANUAL) def test_fail_closed_refusal_names_its_scope_not_the_profile(self): """AC17/AC41: no refusal is profile-wide.""" verdict = mwi.assess_provenance( registry=None, worker_identity=None, env={}, profile="prgs-reviewer", role="reviewer", now=NOW, ) self.assertFalse(verdict["scope"]["profile_wide"]) self.assertEqual(verdict["blocker_kind"], mwi.BLOCKER_NO_ATTACHMENT) class CohortScopingTests(unittest.TestCase): def test_shared_profile_with_distinct_identities_does_not_block(self): """AC40: profile is not a singleton identity.""" registry = _registry() _attach( registry, client="codex", session="s1", generation="g1", profile="prgs-reviewer", ) _attach( registry, client="gemini", session="s2", generation="g2", profile="prgs-reviewer", ) cohort = mwi.classify_cohort(registry.list_workers(), now=NOW) self.assertFalse(cohort["blocked"], cohort["reasons"]) self.assertEqual(cohort["blocker_kind"], mwi.BLOCKER_NONE) self.assertIn("prgs-reviewer", cohort["shared_profiles"]) self.assertTrue(cohort["profile_sharing_permitted"]) self.assertEqual(cohort["blocked_worker_identities"], []) def test_duplicate_cohort_records_block_only_the_offenders(self): registry = _registry() _attach(registry, client="codex", session="s1", generation="gen-contested") _attach(registry, client="gemini", session="s2", generation="gen-contested") _attach(registry, client="claude", session="s3", generation="gen-fine") cohort = mwi.classify_cohort(registry.list_workers(), now=NOW) self.assertTrue(cohort["blocked"]) self.assertEqual(cohort["contested_generations"], ["gen-contested"]) self.assertEqual(len(cohort["blocked_worker_identities"]), 2) def test_mixed_runtime_generations_are_scoped_independently(self): """AC17: one stale generation does not wall unrelated healthy ones.""" registry = _registry() stale = _attach( registry, client="codex", session="s1", generation="gen-stale", ttl=60 ) healthy_a = _attach(registry, client="gemini", session="s2", generation="gen-a") healthy_b = _attach(registry, client="claude", session="s3", generation="gen-b") scoped = mwi.scope_runtime_failure( failure_kind="stale-runtime", worker_identity=stale["identity"], profile="prgs-reviewer", all_live_workers=registry.list_workers(), ) self.assertFalse(scoped["profile_wide"]) self.assertFalse(scoped["fleet_wide"]) self.assertEqual(len(scoped["affected_workers"]), 1) self.assertEqual(scoped["unaffected_worker_count"], 2) unaffected = {w["worker_identity"] for w in scoped["unaffected_workers"]} self.assertEqual(unaffected, {healthy_a["identity"], healthy_b["identity"]}) class HardcodedClientRegressionTests(unittest.TestCase): def test_unknown_client_does_not_resolve_to_codex(self): for name in ("gemini", "antigravity", "grok", "some_new_llm", "", None): with self.subTest(client=name): self.assertNotEqual( mcp_client_reconnect.normalize_client(name), "codex", "an unidentified client must never be handed Codex UI steps", ) def test_known_clients_still_get_their_own_steps(self): self.assertEqual(mcp_client_reconnect.normalize_client("codex"), "codex") self.assertEqual( mcp_client_reconnect.normalize_client("claude_code"), "claude_code" ) def test_generic_steps_do_not_name_a_specific_vendor(self): steps = " ".join(mcp_client_reconnect.operator_ui_steps("gemini")) self.assertNotIn("Codex", steps) def test_reconnect_client_is_derived_from_the_attachment_record(self): registry = _registry() worker = _attach(registry, client="antigravity", session="s", generation="g") verdict = mwi.assess_provenance( registry=registry, worker_identity=worker["identity"], env={}, now=NOW ) self.assertEqual(mwi.reconnect_client_for(verdict), "antigravity") def test_reconnect_client_is_unknown_rather_than_guessed(self): self.assertEqual(mwi.reconnect_client_for({}), mwi.UNKNOWN_CLIENT) class RemoteBindingTests(unittest.TestCase): def test_explicit_prgs_selection_is_honoured(self): resolved = mwi.resolve_bound_remote( requested_remote="prgs", bound_remote="prgs", default_remote="dadeschools" ) self.assertEqual(resolved["remote"], "prgs") self.assertFalse(resolved["drifted"]) def test_omitted_remote_uses_the_binding_not_the_library_default(self): """The reported dadeschools host drift.""" resolved = mwi.resolve_bound_remote( requested_remote=None, bound_remote="prgs", default_remote="dadeschools" ) self.assertEqual(resolved["remote"], "prgs") self.assertNotEqual(resolved["remote"], "dadeschools") self.assertEqual(resolved["resolved_from"], "session_binding") def test_contradicting_the_binding_is_refused(self): resolved = mwi.resolve_bound_remote( requested_remote="dadeschools", bound_remote="prgs", default_remote="dadeschools", ) self.assertEqual(resolved["remote"], "prgs") self.assertTrue(resolved["drifted"]) self.assertFalse(resolved["honoured_request"]) def test_unbound_session_falls_back_and_says_so(self): resolved = mwi.resolve_bound_remote( requested_remote=None, bound_remote=None, default_remote="dadeschools" ) self.assertEqual(resolved["remote"], "dadeschools") self.assertEqual(resolved["resolved_from"], "library_default") self.assertTrue(resolved["reasons"]) class SurfaceAgreementTests(unittest.TestCase): """The reproduced contradiction: two surfaces, one process, two answers.""" def test_namespace_health_and_direct_assessment_agree(self): registry = _registry() worker = _attach( registry, client="gemini", session="sess-agree", generation="gen-agree", profile="prgs-reviewer", ) env = {"GITEA_MCP_PROFILE": "prgs-reviewer", "GITEA_CLIENT_MANAGED": "1"} direct = mwi.assess_provenance( registry=registry, worker_identity=worker["identity"], env=env, profile="prgs-reviewer", ) health = mcp_namespace_health.classify_namespace_probe( "gitea-reviewer", configured=True, registered_tools=["gitea_whoami"], probe_result={"success": True}, probe_source="client_namespace", process={"pid": 4242, "profile": "prgs-reviewer", "env": env}, registry=registry, worker_identity=worker["identity"], ) self.assertEqual(health["provenance"], direct["provenance"]) self.assertEqual(health["is_client_managed"], direct["is_client_managed"]) self.assertEqual(health["worker_identity"], direct["worker_identity"]) self.assertEqual(health["session_id"], "sess-agree") self.assertEqual(health["client_name"], "gemini") def test_namespace_health_can_report_client_managed_at_all(self): """The old derivation was structurally incapable of this.""" env = {"GITEA_CLIENT_MANAGED": "1", "GITEA_MCP_PROFILE": "prgs-author"} health = mcp_namespace_health.classify_namespace_probe( "gitea-author", configured=True, registered_tools=["gitea_whoami"], probe_result={"success": True}, probe_source="client_namespace", process={"pid": 1234, "profile": "prgs-author", "env": env}, ) self.assertTrue( health["is_client_managed"], "a client-managed launch must be reportable as client-managed", ) def test_namespace_health_without_attachment_fails_closed(self): health = mcp_namespace_health.classify_namespace_probe( "gitea-author", configured=True, registered_tools=["gitea_whoami"], probe_result={"success": True}, probe_source="client_namespace", process={"pid": 1234, "profile": "prgs-author", "env": {}}, ) self.assertTrue(health["provenance_fail_closed"]) self.assertEqual(health["provenance"], mwi.PROVENANCE_UNPROVEN) self.assertIsNone(health["session_id"]) def test_no_false_reconnect_loop_for_an_owned_session(self): """A proven owner must not be told to reconnect.""" registry = _registry() worker = _attach(registry, client="claude", session="s", generation="g") verdict = mwi.assess_provenance( registry=registry, worker_identity=worker["identity"], env={}, now=NOW ) self.assertFalse(verdict["fail_closed"]) self.assertEqual(verdict["blocker_kind"], mwi.BLOCKER_NONE) self.assertEqual(verdict["reasons"], []) if __name__ == "__main__": unittest.main()