"""Cross-role allocation handoff consumable by independent workers (#843). Regression coverage for the controller→required-role consume path: * controller allocates author work; independent author adopts successfully * author adoption succeeds after allocating controller process exits * author adoption without sharing controller session identity * wrong-role adoption rejected * concurrent/second adoption rejected without state corruption * terminal allocation adoption rejected * successful adoption produces authoritative ownership evidence * genuine abandoned-lease recovery remains valid * process_work_queue / allocate results include consume identifiers * same-role allocation behavior remains compatible """ from __future__ import annotations import os import tempfile import unittest from datetime import timedelta from unittest.mock import patch from allocator_service import ( ALLOCATION_MODE_CROSS_ROLE, ALLOCATION_MODE_ROLE_SCOPED, OUTCOME_ASSIGNED, ROLE_AUTHOR, ROLE_CONTROLLER, ROLE_REVIEWER, WorkCandidate, allocate_next_work, ) from control_plane_db import ControlPlaneDB, ForeignLeaseError, _ts, _utc_now import lease_lifecycle as ll class CrossRoleHandoffTest(unittest.TestCase): def setUp(self) -> None: self._tmp = tempfile.TemporaryDirectory() self.db_path = os.path.join(self._tmp.name, "cp.sqlite3") self.db = ControlPlaneDB(self.db_path) self.db.upsert_session( session_id="ctrl-session", role="controller", profile="prgs-controller", pid=99999999, # dead-looking pid ) self.db.upsert_session( session_id="author-worker", role="author", profile="prgs-author", pid=os.getpid(), ) self.db.upsert_session( session_id="author-worker-2", role="author", profile="prgs-author", pid=os.getpid(), ) self.db.upsert_session( session_id="reviewer-worker", role="reviewer", profile="prgs-reviewer", pid=os.getpid(), ) self.wt = self._tmp.name def tearDown(self) -> None: self._tmp.cleanup() def _ready_issue(self, number: int = 843, title: str = "handoff target") -> WorkCandidate: return WorkCandidate( kind="issue", number=number, labels=("status:ready", "type:bug"), title=title, priority=20, ) def _controller_allocate(self, number: int = 843, **kwargs): defaults = dict( db=self.db, session_id="ctrl-session", role=ROLE_CONTROLLER, remote="prgs", org="Scaled-Tech-Consulting", repo="Gitea-Tools", candidates=[self._ready_issue(number)], apply=True, profile_name="prgs-controller", username="controller-user", allocation_mode=ALLOCATION_MODE_CROSS_ROLE, ) defaults.update(kwargs) return allocate_next_work(**defaults) def test_controller_allocates_author_independent_author_adopts(self) -> None: res = self._controller_allocate() self.assertEqual(res["outcome"], OUTCOME_ASSIGNED) self.assertEqual(res["required_role"], ROLE_AUTHOR) self.assertIn("consume_allocation", res) consume = res["consume_allocation"] self.assertEqual(consume["tool"], "gitea_adopt_workflow_lease") self.assertEqual(consume["required_role"], ROLE_AUTHOR) self.assertFalse(consume["controller_session_required"]) lid = res["assignment"]["lease_id"] self.assertEqual(consume["lease_id"], lid) self.assertIn(lid, res["next_valid_command"]) adopted = ll.adopt_lease( self.db, lease_id=lid, adopter_session_id="author-worker", role=ROLE_AUTHOR, worktree_path=self.wt, ) self.assertTrue(adopted["success"]) self.assertEqual(adopted["outcome"], "adopted_cross_role_handoff") self.assertEqual(adopted["adopted_by_session_id"], "author-worker") self.assertEqual(adopted["adopted_from_session_id"], "ctrl-session") raw = adopted["read_after_write"] self.assertEqual(raw["session_id"], "author-worker") self.assertEqual(raw["adopted_by_session_id"], "author-worker") self.assertEqual(raw["status"], "active") self.assertEqual(raw["phase"], "adopted") # Authoritative re-read state = self.db.get_lease_workflow_state(lid) self.assertEqual(state["lease"]["session_id"], "author-worker") self.assertEqual(state["lease"]["adopted_by_session_id"], "author-worker") self.assertEqual(state["assignment"]["session_id"], "author-worker") self.assertEqual(state["provenance"]["handoff_status"], "adopted") def test_author_adoption_after_controller_process_exits(self) -> None: res = self._controller_allocate(number=900) lid = res["assignment"]["lease_id"] # Force owner_pid dead + freshness stale_dead_process import sqlite3 conn = sqlite3.connect(self.db_path) try: conn.execute( "UPDATE leases SET owner_pid = 99999999 WHERE lease_id = ?", (lid,), ) conn.commit() finally: conn.close() state = self.db.get_lease_workflow_state(lid) fr = ll.classify_lease_freshness( state["lease"], pid_checker=lambda _p: False ) self.assertEqual(fr["freshness"], "stale_dead_process") adopted = ll.adopt_lease( self.db, lease_id=lid, adopter_session_id="author-worker", role=ROLE_AUTHOR, worktree_path=self.wt, ) self.assertEqual(adopted["outcome"], "adopted_cross_role_handoff") self.assertEqual(adopted["adopted_by_session_id"], "author-worker") # No abandon required state2 = self.db.get_lease_workflow_state(lid) self.assertEqual(state2["lease"]["status"], "active") self.assertNotEqual(state2["lease"]["status"], "abandoned") def test_adoption_without_sharing_controller_session_identity(self) -> None: res = self._controller_allocate(number=901) lid = res["assignment"]["lease_id"] adopted = ll.adopt_lease( self.db, lease_id=lid, adopter_session_id="author-worker", role=ROLE_AUTHOR, worktree_path=self.wt, ) self.assertNotEqual(adopted["adopted_by_session_id"], "ctrl-session") self.assertFalse(adopted["same_owner"]) self.assertEqual(adopted["adopted_from_session_id"], "ctrl-session") def test_wrong_role_adoption_rejected(self) -> None: res = self._controller_allocate(number=902) lid = res["assignment"]["lease_id"] with self.assertRaises(ll.LeaseLifecycleError) as ctx: ll.adopt_lease( self.db, lease_id=lid, adopter_session_id="reviewer-worker", role=ROLE_REVIEWER, ) self.assertIn("wrong role", str(ctx.exception).lower()) # State unchanged state = self.db.get_lease_workflow_state(lid) self.assertEqual(state["lease"]["session_id"], "ctrl-session") self.assertIsNone(state["lease"].get("adopted_by_session_id") or None) self.assertEqual(state["provenance"]["handoff_status"], "pending") def test_second_adoption_rejected_without_corruption(self) -> None: res = self._controller_allocate(number=903) lid = res["assignment"]["lease_id"] first = ll.adopt_lease( self.db, lease_id=lid, adopter_session_id="author-worker", role=ROLE_AUTHOR, worktree_path=self.wt, ) self.assertEqual(first["outcome"], "adopted_cross_role_handoff") with self.assertRaises(ll.LeaseLifecycleError): ll.adopt_lease( self.db, lease_id=lid, adopter_session_id="author-worker-2", role=ROLE_AUTHOR, worktree_path=self.wt, ) state = self.db.get_lease_workflow_state(lid) self.assertEqual(state["lease"]["session_id"], "author-worker") self.assertEqual(state["lease"]["adopted_by_session_id"], "author-worker") self.assertEqual(state["assignment"]["session_id"], "author-worker") self.assertEqual(state["lease"]["status"], "active") def test_terminal_allocation_adoption_rejected(self) -> None: res = self._controller_allocate(number=904) lid = res["assignment"]["lease_id"] # Abandon as terminal proof = ll.AbandonProof( dead_process=True, missing_worktree=True, no_open_pr=True, no_live_mutation_risk=True, owner_pid=99999999, worktree_path="/nonexistent/for-843", ) # Attach dead pid / missing wt for abandon eligibility import sqlite3 conn = sqlite3.connect(self.db_path) try: conn.execute( "UPDATE leases SET owner_pid = 99999999, worktree_path = ? WHERE lease_id = ?", ("/nonexistent/for-843", lid), ) conn.commit() finally: conn.close() abandoned = ll.abandon_lease( self.db, lease_id=lid, requester_session_id="author-worker", proof=proof, ) self.assertEqual(abandoned["outcome"], "abandoned") with self.assertRaises(ll.LeaseLifecycleError) as ctx: ll.adopt_lease( self.db, lease_id=lid, adopter_session_id="author-worker", role=ROLE_AUTHOR, ) self.assertIn("abandoned", str(ctx.exception).lower()) def test_successful_adoption_read_after_write_ownership(self) -> None: res = self._controller_allocate(number=905) lid = res["assignment"]["lease_id"] adopted = ll.adopt_lease( self.db, lease_id=lid, adopter_session_id="author-worker", role=ROLE_AUTHOR, worktree_path=self.wt, ) raw = adopted["read_after_write"] self.assertEqual(raw["lease_id"], lid) self.assertEqual(raw["session_id"], "author-worker") self.assertEqual(raw["adopted_by_session_id"], "author-worker") self.assertEqual(raw["adopted_from_session_id"], "ctrl-session") # Re-fetch proves durable write state = self.db.get_lease_workflow_state(lid) self.assertEqual(state["lease"]["session_id"], raw["session_id"]) self.assertEqual( state["lease"]["adopted_by_session_id"], raw["adopted_by_session_id"] ) def test_genuine_abandoned_recovery_still_valid(self) -> None: """Same-role author lease abandoned remains reclaimable via abandon path.""" same = allocate_next_work( self.db, session_id="author-worker", role=ROLE_AUTHOR, remote="prgs", org="Scaled-Tech-Consulting", repo="Gitea-Tools", candidates=[self._ready_issue(906, "same-role")], apply=True, profile_name="prgs-author", username="author-user", allocation_mode=ALLOCATION_MODE_ROLE_SCOPED, ) self.assertEqual(same["outcome"], OUTCOME_ASSIGNED) lid = same["assignment"]["lease_id"] import sqlite3 conn = sqlite3.connect(self.db_path) try: conn.execute( "UPDATE leases SET owner_pid = 99999999, worktree_path = ? WHERE lease_id = ?", ("/nonexistent/same-role", lid), ) conn.commit() finally: conn.close() proof = ll.AbandonProof( dead_process=True, missing_worktree=True, no_open_pr=True, no_live_mutation_risk=True, owner_pid=99999999, worktree_path="/nonexistent/same-role", ) abandoned = ll.abandon_lease( self.db, lease_id=lid, requester_session_id="author-worker-2", proof=proof, ) self.assertEqual(abandoned["outcome"], "abandoned") # Foreign author cannot handoff-consume an abandoned non-handoff lease with self.assertRaises(ll.LeaseLifecycleError): ll.adopt_lease( self.db, lease_id=lid, adopter_session_id="author-worker-2", role=ROLE_AUTHOR, ) # Reclaim path still works for expired/abandoned after force-expire reclaimed = ll.reclaim_expired_lease( self.db, lease_id=lid, session_id="author-worker-2", role=ROLE_AUTHOR, worktree_path=self.wt, ) self.assertEqual(reclaimed["outcome"], "reclaimed") self.assertEqual(reclaimed["assignment"]["session_id"], "author-worker-2") def test_allocate_payload_includes_consume_identifiers(self) -> None: res = self._controller_allocate(number=907) self.assertIn("consume_allocation", res) c = res["consume_allocation"] for key in ( "tool", "lease_id", "assignment_id", "required_role", "required_profile", "required_namespace", "instructions", "handoff_status", ): self.assertIn(key, c) self.assertEqual(c["required_namespace"], "gitea-author") self.assertEqual(c["required_profile"], "prgs-author") self.assertIn("gitea_adopt_workflow_lease", c["instructions"]) self.assertTrue(res["lease_proof"]["cross_role_handoff"]) self.assertEqual(res["lease_proof"]["handoff_status"], "pending") def test_same_role_allocation_remains_compatible(self) -> None: res = allocate_next_work( self.db, session_id="author-worker", role=ROLE_AUTHOR, remote="prgs", org="Scaled-Tech-Consulting", repo="Gitea-Tools", candidates=[self._ready_issue(908)], apply=True, profile_name="prgs-author", username="author-user", ) self.assertEqual(res["outcome"], OUTCOME_ASSIGNED) self.assertNotIn("consume_allocation", res) lid = res["assignment"]["lease_id"] state = self.db.get_lease_workflow_state(lid) # No cross-role handoff provenance prov = state.get("provenance") or {} self.assertFalse(prov.get("cross_role_handoff")) # Owner resume still works resume = ll.adopt_lease( self.db, lease_id=lid, adopter_session_id="author-worker", role=ROLE_AUTHOR, worktree_path=self.wt, ) self.assertTrue(resume["same_owner"]) self.assertEqual(resume["outcome"], "adopted_owner_resume") def test_inspect_points_required_role_at_consume(self) -> None: res = self._controller_allocate(number=909) lid = res["assignment"]["lease_id"] decision = ll.inspect_lease( self.db, lid, caller_session_id="author-worker" ) self.assertEqual( decision["safe_next_action"], ll.SAFE_CONSUME_CROSS_ROLE ) self.assertFalse(decision["block"]) self.assertEqual(decision["required_role"], ROLE_AUTHOR) def test_db_cas_rejects_concurrent_second_consume(self) -> None: res = self._controller_allocate(number=910) lid = res["assignment"]["lease_id"] # First consume via DB layer directly first = self.db.adopt_lease( lease_id=lid, adopter_session_id="author-worker", role=ROLE_AUTHOR, worktree_path=self.wt, provenance={ "cross_role_handoff": True, "handoff_status": "adopted", "required_role": "author", }, ) self.assertEqual(first["outcome"], "adopted_cross_role_handoff") # Second CAS must fail with self.assertRaises(ForeignLeaseError): self.db.adopt_lease( lease_id=lid, adopter_session_id="author-worker-2", role=ROLE_AUTHOR, worktree_path=self.wt, provenance={ "cross_role_handoff": True, "handoff_status": "pending", "required_role": "author", }, ) state = self.db.get_lease_workflow_state(lid) self.assertEqual(state["lease"]["session_id"], "author-worker") if __name__ == "__main__": unittest.main()