Compare commits

..
Author SHA1 Message Date
sysadminandClaude Opus 4.8 7f97de1ed6 feat: gate author work on early duplicate-work detection (Closes #400)
Add author_duplicate_work_gate and enforce it at claim, lock, and PR
creation. Expose gitea_assess_author_duplicate_work for pre-commit/push
checks and extend work-issue final-report verification.

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-07 13:08:10 -04:00
10 changed files with 485 additions and 1018 deletions
+198
View File
@@ -0,0 +1,198 @@
"""Early duplicate-work detection for author work-issue flows (#400)."""
from __future__ import annotations
import re
from typing import Any
from issue_claim_heartbeat import (
_linked_open_pr,
_matching_branch_names,
classify_issue_claim,
)
STAGES = (
"claim",
"lock",
"worktree",
"edit",
"commit",
"push",
"create_pr",
)
ELIGIBILITY_OPEN_PR_EXISTS = "OPEN_PR_EXISTS"
ELIGIBILITY_DUPLICATE_BRANCH_EXISTS = "DUPLICATE_BRANCH_EXISTS"
ELIGIBILITY_ACTIVE_CLAIM = "ACTIVE_CLAIM_BY_OTHER"
ELIGIBILITY_CLEAR = "CLEAR"
def assess_author_duplicate_work(
issue_number: int,
*,
stage: str,
open_prs: list[dict] | None = None,
branch_names: list[str] | None = None,
claim_entry: dict | None = None,
matching_branches: list[str] | None = None,
allow_stale_takeover: bool = False,
) -> dict[str, Any]:
"""Fail closed when duplicate work is detected before author mutations (#400)."""
stage_norm = (stage or "").strip().lower()
if stage_norm not in STAGES:
return {
"allowed": False,
"block": True,
"eligibility_class": "INVALID_STAGE",
"stage": stage_norm or None,
"reasons": [f"unknown duplicate-work stage {stage!r}"],
"safe_next_action": f"use one of: {', '.join(STAGES)}",
}
prs = list(open_prs or [])
linked_pr = _linked_open_pr(int(issue_number), prs)
branches = list(
matching_branches
if matching_branches is not None
else _matching_branch_names(int(issue_number), list(branch_names or []))
)
reasons: list[str] = []
eligibility = ELIGIBILITY_CLEAR
if linked_pr:
eligibility = ELIGIBILITY_OPEN_PR_EXISTS
reasons.append(
f"open PR #{linked_pr.get('number')} already covers issue "
f"#{issue_number}"
)
if branches and stage_norm in {"claim", "lock", "worktree", "edit"}:
if eligibility == ELIGIBILITY_CLEAR:
eligibility = ELIGIBILITY_DUPLICATE_BRANCH_EXISTS
reasons.append(
f"remote branch(es) already exist for issue #{issue_number}: "
f"{', '.join(branches)}"
)
entry = claim_entry or {}
claim_status = (entry.get("status") or "").strip()
if claim_status in {"active", "awaiting_review"} and stage_norm == "claim":
if entry.get("reclaimable") and allow_stale_takeover:
pass
elif claim_status == "active" and not entry.get("reclaimable"):
if eligibility == ELIGIBILITY_CLEAR:
eligibility = ELIGIBILITY_ACTIVE_CLAIM
reasons.append(
f"issue #{issue_number} has active claim "
f"(status={claim_status})"
)
elif claim_status == "awaiting_review" and stage_norm == "claim":
if not linked_pr:
reasons.append(
f"issue #{issue_number} is awaiting_review but no linked open PR "
"was supplied for duplicate-work proof"
)
allowed = not reasons
outcome = "duplicate_work_prevented" if not allowed else "clear"
if stage_norm == "create_pr" and not allowed:
outcome = "duplicate_pr_prevented"
elif stage_norm in {"commit", "push"} and not allowed:
outcome = "duplicate_push_prevented" if stage_norm == "push" else "duplicate_commit_prevented"
return {
"allowed": allowed,
"block": not allowed,
"eligibility_class": eligibility if not allowed else ELIGIBILITY_CLEAR,
"stage": stage_norm,
"linked_open_pr": linked_pr.get("number") if linked_pr else None,
"matching_branches": branches,
"claim_status": claim_status or None,
"outcome": outcome,
"reasons": reasons,
"safe_next_action": (
"stop without edits/commit/push/PR; produce reconciliation handoff "
"preserving local work only"
if not allowed
else "proceed"
),
}
def build_claim_entry_from_classification(classification: dict) -> dict:
"""Map ``classify_issue_claim`` output to gate claim metadata."""
return {
"status": classification.get("status"),
"reclaimable": classification.get("reclaimable"),
"linked_open_pr": classification.get("linked_open_pr"),
}
_DUPLICATE_OUTCOME_RE = re.compile(
r"(duplicate\s+(?:pr|branch|commit|push)\s+prevented|"
r"duplicate\s+work\s+not\s+prevented|reconciliation\s+handoff)",
re.IGNORECASE,
)
def assess_work_issue_duplicate_prevention_report(report_text: str) -> dict:
"""#400: work-issue reports must state duplicate-work prevention outcome."""
text = report_text or ""
if "duplicate work" not in text.lower() and "duplicate pr" not in text.lower():
return {
"proven": True,
"block": False,
"reasons": [],
"safe_next_action": "proceed",
}
if _DUPLICATE_OUTCOME_RE.search(text):
return {
"proven": True,
"block": False,
"reasons": [],
"safe_next_action": "proceed",
}
return {
"proven": False,
"block": True,
"reasons": [
"duplicate-work discussion must name a prevention outcome "
"(duplicate PR/branch/commit/push prevented, or duplicate work "
"not prevented, or reconciliation handoff)"
],
"safe_next_action": "state exact duplicate-work prevention class in final report",
}
def classify_and_assess(
issue: dict,
*,
stage: str,
comments: list[dict] | None = None,
open_prs: list[dict] | None = None,
branch_names: list[str] | None = None,
allow_stale_takeover: bool = False,
) -> dict[str, Any]:
"""Combine claim classification with duplicate-work gate assessment."""
issue_number = int(issue.get("number") or 0)
claim = classify_issue_claim(
issue=issue,
comments=comments or [],
open_prs=open_prs or [],
branch_names=branch_names or [],
)
assessment = assess_author_duplicate_work(
issue_number,
stage=stage,
open_prs=open_prs,
branch_names=branch_names,
claim_entry=build_claim_entry_from_classification(claim),
matching_branches=claim.get("matching_branches"),
allow_stale_takeover=allow_stale_takeover,
)
return {
"issue_number": issue_number,
"claim": claim,
"duplicate_work": assessment,
}
-41
View File
@@ -490,45 +490,6 @@ def _rule_reviewer_validation_failure_history(
] ]
def _rule_reviewer_stale_head_proof(report_text: str) -> list[dict[str, str]]:
from pr_work_lease import assess_reviewer_stale_head_final_report
result = assess_reviewer_stale_head_final_report(report_text)
if result.get("proven"):
return []
return _findings_from_reasons(
"reviewer.stale_head_proof",
result.get("reasons") or [],
field="Stale-head proof",
severity="block",
safe_next_action=(
"state reviewed head SHA, live head before approval/merge, and "
"whether any push occurred during validation"
),
)
def _rule_conflict_fix_push_proof(report_text: str) -> list[dict[str, str]]:
from pr_work_lease import assess_conflict_fix_final_report
text = report_text or ""
if "conflict-fix" not in text.lower() and "conflict fix" not in text.lower():
return []
result = assess_conflict_fix_final_report(text)
if result.get("proven"):
return []
return _findings_from_reasons(
"author.conflict_fix_push_proof",
result.get("reasons") or [],
field="Conflict-fix push proof",
severity="block",
safe_next_action=(
"state branch head before/after push, reviewer lease status, "
"fast-forward status, and whether any reviewer was active"
),
)
def _rule_reviewer_validation_command(report_text: str) -> list[dict[str, str]]: def _rule_reviewer_validation_command(report_text: str) -> list[dict[str, str]]:
text = report_text or "" text = report_text or ""
if not _BARE_PYTEST_RE.search(text): if not _BARE_PYTEST_RE.search(text):
@@ -948,7 +909,6 @@ _RULES_BY_TASK: dict[str, list[Callable[..., list[dict[str, str]]]]] = {
_rule_reviewer_target_branch_freshness, _rule_reviewer_target_branch_freshness,
_rule_reviewer_mutation_ledger, _rule_reviewer_mutation_ledger,
_rule_reviewer_review_mutation, _rule_reviewer_review_mutation,
_rule_reviewer_stale_head_proof,
], ],
"reconcile_already_landed": [ "reconcile_already_landed": [
_rule_reconcile_controller_handoff, _rule_reconcile_controller_handoff,
@@ -970,7 +930,6 @@ _RULES_BY_TASK: dict[str, list[Callable[..., list[dict[str, str]]]]] = {
_rule_shared_controller_handoff, _rule_shared_controller_handoff,
_rule_shared_email_disclosure, _rule_shared_email_disclosure,
_rule_reviewer_vague_mutations_none, _rule_reviewer_vague_mutations_none,
_rule_conflict_fix_push_proof,
], ],
"issue_filing": [ "issue_filing": [
_rule_shared_controller_handoff, _rule_shared_controller_handoff,
+142 -246
View File
@@ -503,11 +503,11 @@ import issue_lock_worktree # noqa: E402
import already_landed_reconcile # noqa: E402 import already_landed_reconcile # noqa: E402
import author_mutation_worktree # noqa: E402 import author_mutation_worktree # noqa: E402
import issue_claim_heartbeat # noqa: E402 import issue_claim_heartbeat # noqa: E402
import author_duplicate_work_gate # noqa: E402
import merged_cleanup_reconcile # noqa: E402 import merged_cleanup_reconcile # noqa: E402
import reconciler_profile # noqa: E402 import reconciler_profile # noqa: E402
import reconciliation_workflow # noqa: E402 import reconciliation_workflow # noqa: E402
import review_merge_state_machine # noqa: E402 import review_merge_state_machine # noqa: E402
import pr_work_lease # noqa: E402
import native_mcp_preference # noqa: E402 import native_mcp_preference # noqa: E402
@@ -1050,6 +1050,76 @@ def gitea_create_issue(
return _with_optional_url({"number": data["number"]}, data.get("html_url")) return _with_optional_url({"number": data["number"]}, data.get("html_url"))
def _list_repo_branch_names(h: str, o: str, r: str, auth: str, *, limit: int = 200) -> list[str]:
branches = api_get_all(f"{repo_api_url(h, o, r)}/branches", auth, limit=limit)
return [_branch_entry_name(branch) for branch in branches]
def _gather_author_duplicate_work_context(
issue_number: int,
*,
h: str,
o: str,
r: str,
auth: str,
exclude_branch_name: str | None = None,
) -> dict:
base = repo_api_url(h, o, r)
issue = api_request("GET", f"{base}/issues/{issue_number}", auth)
comments = api_request("GET", f"{base}/issues/{issue_number}/comments", auth) or []
open_prs = api_get_all(f"{base}/pulls?state=open", auth)
branch_names = _list_repo_branch_names(h, o, r, auth)
if exclude_branch_name:
branch_names = [
name for name in branch_names
if name != exclude_branch_name
]
return {
"issue": issue,
"comments": comments,
"open_prs": open_prs,
"branch_names": branch_names,
}
def _enforce_author_duplicate_work_gate(
issue_number: int,
stage: str,
*,
h: str,
o: str,
r: str,
auth: str,
allow_stale_takeover: bool = False,
exclude_branch_name: str | None = None,
) -> dict:
"""Fail closed when duplicate work is detected (#400)."""
ctx = _gather_author_duplicate_work_context(
issue_number,
h=h,
o=o,
r=r,
auth=auth,
exclude_branch_name=exclude_branch_name,
)
result = author_duplicate_work_gate.classify_and_assess(
ctx["issue"],
stage=stage,
comments=ctx["comments"],
open_prs=ctx["open_prs"],
branch_names=ctx["branch_names"],
allow_stale_takeover=allow_stale_takeover,
)
assessment = result.get("duplicate_work") or {}
if assessment.get("block"):
reasons = "; ".join(assessment.get("reasons") or ["duplicate work detected"])
raise RuntimeError(
f"Author duplicate-work gate (#400) blocked at stage '{stage}': "
f"{reasons} (fail closed)"
)
return result
@mcp.tool() @mcp.tool()
def gitea_lock_issue( def gitea_lock_issue(
issue_number: int, issue_number: int,
@@ -1112,48 +1182,17 @@ def gitea_lock_issue(
issue_lock_worktree.format_issue_lock_worktree_error(lock_assessment) issue_lock_worktree.format_issue_lock_worktree_error(lock_assessment)
) )
# 2. Check if the issue already has an open PR (reuse protection)
h, o, r = _resolve(remote, host, org, repo) h, o, r = _resolve(remote, host, org, repo)
auth = _auth(h) auth = _auth(h)
url = f"{repo_api_url(h, o, r)}/pulls?state=open" _enforce_author_duplicate_work_gate(
issue_number,
try: "lock",
prs = api_get_all(url, auth) h=h,
except Exception as e: o=o,
raise RuntimeError(f"Could not list open PRs to verify issue lock: {e}") r=r,
auth=auth,
for pr in prs: exclude_branch_name=branch_name,
pr_head = pr.get("head", {}).get("ref", "") )
pr_title = pr.get("title", "")
pr_body = pr.get("body", "")
if expected_pattern in pr_head:
raise ValueError(
f"Issue #{issue_number} is already tied to an open PR (PR #{pr.get('number')}, branch '{pr_head}') (fail closed)"
)
patterns = [
f"closes #{issue_number}",
f"fixes #{issue_number}",
]
text_to_check = f"{pr_title} {pr_body}".lower()
if any(p in text_to_check for p in patterns):
raise ValueError(
f"Issue #{issue_number} is already tied to an open PR (PR #{pr.get('number')}) via Closes/Fixes reference (fail closed)"
)
branch_url = f"{repo_api_url(h, o, r)}/branches"
try:
branches = api_get_all(branch_url, auth)
except Exception as e:
raise RuntimeError(f"Could not list branches to verify issue lock: {e}")
for branch in branches:
name = _branch_entry_name(branch)
if expected_pattern in name:
raise ValueError(
f"Issue #{issue_number} already has matching branch '{name}' "
"(fail closed)"
)
work_lease = _build_author_issue_work_lease( work_lease = _build_author_issue_work_lease(
issue_number=issue_number, issue_number=issue_number,
@@ -1274,6 +1313,17 @@ def gitea_create_pr(
f"PR head branch '{head}' does not match locked branch '{locked_branch}' (fail closed)" f"PR head branch '{head}' does not match locked branch '{locked_branch}' (fail closed)"
) )
auth = _auth(h)
_enforce_author_duplicate_work_gate(
int(locked_issue),
"create_pr",
h=h,
o=o,
r=r,
auth=auth,
exclude_branch_name=locked_branch,
)
# Check for forbidden terms anywhere in title/body # Check for forbidden terms anywhere in title/body
forbidden_terms = ["equivalent", "related", "same as"] forbidden_terms = ["equivalent", "related", "same as"]
text_to_check = f"{title} {body}".lower() text_to_check = f"{title} {body}".lower()
@@ -1290,7 +1340,6 @@ def gitea_create_pr(
f"PR title or body must contain 'Closes #{locked_issue}' or 'Fixes #{locked_issue}' exactly to ensure durable tracking (fail closed)" f"PR title or body must contain 'Closes #{locked_issue}' or 'Fixes #{locked_issue}' exactly to ensure durable tracking (fail closed)"
) )
auth = _auth(h)
url = f"{repo_api_url(h, o, r)}/pulls" url = f"{repo_api_url(h, o, r)}/pulls"
payload = {"title": title, "body": body, "head": head, "base": base} payload = {"title": title, "body": body, "head": head, "base": base}
meta = {"title": title, "head": head, "base": base} meta = {"title": title, "head": head, "base": base}
@@ -2198,50 +2247,6 @@ def gitea_get_pr_review_feedback(
} }
def _list_pr_lease_comments(
pr_number: int,
*,
remote: str,
host: str | None,
org: str | None,
repo: str | None,
limit: int = 100,
) -> list[dict]:
"""Fetch PR/issue thread comments used for reviewer/conflict-fix leases."""
h, o, r = _resolve(remote, host, org, repo)
auth = _auth(h)
api = f"{repo_api_url(h, o, r)}/issues/{pr_number}/comments"
comments = api_request("GET", api, auth) or []
return list(comments[:limit])
def _pr_work_lease_reviewer_block(
*,
pr_number: int,
reviewed_head_sha: str | None,
live_head_sha: str | None,
mutation: str,
remote: str,
host: str | None,
org: str | None,
repo: str | None,
) -> dict:
comments = _list_pr_lease_comments(
pr_number,
remote=remote,
host=host,
org=org,
repo=repo,
)
return pr_work_lease.assess_reviewer_mutation_blocked(
pr_number=pr_number,
comments=comments,
reviewed_head_sha=reviewed_head_sha,
live_head_sha=live_head_sha,
mutation=mutation,
)
def _evaluate_pr_review_submission( def _evaluate_pr_review_submission(
pr_number: int, pr_number: int,
action: str, action: str,
@@ -2326,11 +2331,6 @@ def _evaluate_pr_review_submission(
lock = _load_review_decision_lock() or {} lock = _load_review_decision_lock() or {}
if live and lock.get("ready_expected_head_sha"): if live and lock.get("ready_expected_head_sha"):
pinned_sha = lock.get("ready_expected_head_sha") pinned_sha = lock.get("ready_expected_head_sha")
if live and not pinned_sha:
reasons.append(
"reviewed head SHA required before live review mutation (fail closed, #399)"
)
return result
if pinned_sha and actual_sha and pinned_sha != actual_sha: if pinned_sha and actual_sha and pinned_sha != actual_sha:
reasons.append( reasons.append(
"expected head SHA does not match current PR head (fail closed)" "expected head SHA does not match current PR head (fail closed)"
@@ -2340,21 +2340,6 @@ def _evaluate_pr_review_submission(
reasons.append("PR head SHA unavailable (fail closed)") reasons.append("PR head SHA unavailable (fail closed)")
return result return result
lease_block = _pr_work_lease_reviewer_block(
pr_number=pr_number,
reviewed_head_sha=pinned_sha,
live_head_sha=actual_sha,
mutation=action,
remote=remote,
host=host,
org=org,
repo=repo,
)
if lease_block.get("block"):
reasons.extend(lease_block.get("reasons") or [])
result["pr_work_lease"] = lease_block
return result
result["would_perform"] = True result["would_perform"] = True
if not live: if not live:
reasons.append( reasons.append(
@@ -2497,39 +2482,6 @@ def gitea_mark_final_review_decision(
f"{sorted(_REVIEW_ACTIONS)}" f"{sorted(_REVIEW_ACTIONS)}"
], ],
} }
if not (expected_head_sha or "").strip():
return {
"marked_ready": False,
"reasons": [
"expected_head_sha required before marking final review "
"decision (fail closed, #399)"
],
}
elig = gitea_check_pr_eligibility(
pr_number=pr_number,
action="review",
remote=remote,
host=None,
org=org,
repo=repo,
)
live_head = elig.get("head_sha")
lease_block = _pr_work_lease_reviewer_block(
pr_number=pr_number,
reviewed_head_sha=expected_head_sha,
live_head_sha=live_head,
mutation="mark_ready",
remote=remote,
host=None,
org=org,
repo=repo,
)
if lease_block.get("block"):
return {
"marked_ready": False,
"reasons": lease_block.get("reasons") or [],
"pr_work_lease": lease_block,
}
if action == "request_changes": if action == "request_changes":
# Duplicate request-changes suppression (#332): an unresolved # Duplicate request-changes suppression (#332): an unresolved
# REQUEST_CHANGES at the current head must not be duplicated. # REQUEST_CHANGES at the current head must not be duplicated.
@@ -3153,27 +3105,8 @@ def gitea_merge_pr(
result["permission_report"] = elig["permission_report"] result["permission_report"] = elig["permission_report"]
return result return result
# Gate 4 — reviewed head SHA is mandatory and must match live PR head (#399). # Gate 4 — head SHA must match if the caller pinned a reviewed SHA.
actual_sha = result["head_sha"] actual_sha = result["head_sha"]
if not (expected_head_sha or "").strip():
reasons.append(
"expected_head_sha required before merge (fail closed, #399)"
)
return result
lease_block = _pr_work_lease_reviewer_block(
pr_number=pr_number,
reviewed_head_sha=expected_head_sha,
live_head_sha=actual_sha,
mutation="merge",
remote=remote,
host=host,
org=org,
repo=repo,
)
if lease_block.get("block"):
reasons.extend(lease_block.get("reasons") or [])
result["pr_work_lease"] = lease_block
return result
if expected_head_sha and actual_sha and expected_head_sha != actual_sha: if expected_head_sha and actual_sha and expected_head_sha != actual_sha:
reasons.append( reasons.append(
"expected head SHA does not match current PR head (fail closed)" "expected head SHA does not match current PR head (fail closed)"
@@ -6120,6 +6053,14 @@ def gitea_mark_issue(
) )
if action == "start": if action == "start":
_enforce_author_duplicate_work_gate(
issue_number,
"claim",
h=h,
o=o,
r=r,
auth=auth,
)
with _audited("label_issue", host=h, remote=remote, org=o, repo=r, with _audited("label_issue", host=h, remote=remote, org=o, repo=r,
issue_number=issue_number, issue_number=issue_number,
request_metadata={"op": "add", "label": "status:in-progress"}): request_metadata={"op": "add", "label": "status:in-progress"}):
@@ -6202,107 +6143,62 @@ def gitea_post_heartbeat(
@mcp.tool() @mcp.tool()
def gitea_acquire_conflict_fix_lease( def gitea_assess_author_duplicate_work(
pr_number: int, issue_number: int,
branch: str, stage: str,
worktree_path: str,
head_before: str,
remote: str = "dadeschools", remote: str = "dadeschools",
host: str | None = None, host: str | None = None,
org: str | None = None, org: str | None = None,
repo: str | None = None, repo: str | None = None,
allow_stale_takeover: bool = False,
exclude_branch_name: str | None = None,
) -> dict: ) -> dict:
"""Acquire a conflict-fix lease on a PR branch before pushing (#399).""" """Read-only: assess duplicate-work risk before author mutations (#400).
blocked = _profile_permission_block(
task_capability_map.required_permission("comment_issue"))
if blocked:
return blocked
verify_preflight_purity(remote, worktree_path=worktree_path)
comments = _list_pr_lease_comments(
pr_number,
remote=remote,
host=host,
org=org,
repo=repo,
)
reviewer_lease = pr_work_lease.find_active_reviewer_lease(
comments, pr_number=pr_number)
if reviewer_lease:
return {
"acquired": False,
"reasons": [
f"active reviewer lease on PR #{pr_number}; cannot acquire "
"conflict-fix lease (fail closed)"
],
"active_reviewer_lease": reviewer_lease,
}
profile_name = get_profile().get("profile_name") or "unknown"
body = pr_work_lease.format_conflict_fix_lease_body(
pr_number=pr_number,
branch=branch,
worktree=worktree_path,
profile=profile_name,
head_before=head_before,
reviewer_active=bool(reviewer_lease),
)
posted = _post_structured_issue_comment(
issue_number=pr_number,
body=body,
remote=remote,
host=host,
org=org,
repo=repo,
audit_op="conflict_fix_lease_acquire",
)
return {
"acquired": posted.get("success", False),
"pr_number": pr_number,
"branch": branch,
"worktree_path": worktree_path,
"head_before": head_before,
"comment_id": posted.get("comment_id"),
"active_reviewer_lease": reviewer_lease,
"reasons": [] if posted.get("success") else ["lease comment post failed"],
}
Call at claim, lock, worktree, edit, commit, push, and create_pr stages.
@mcp.tool() """
def gitea_assess_conflict_fix_push(
pr_number: int,
branch_head_before: str,
branch_head_after: str,
worktree_path: str,
push_cwd: str,
is_fast_forward: bool = True,
remote: str = "dadeschools",
host: str | None = None,
org: str | None = None,
repo: str | None = None,
) -> dict:
"""Read-only pre-push gate for author conflict-fix sessions (#399)."""
read_block = _profile_operation_gate("gitea.read") read_block = _profile_operation_gate("gitea.read")
if read_block: if read_block:
return { return {
"push_allowed": False, "success": False,
"reasons": read_block, "reasons": read_block,
"permission_report": _permission_block_report("gitea.read"), "permission_report": _permission_block_report("gitea.read"),
} }
comments = _list_pr_lease_comments( h, o, r = _resolve(remote, host, org, repo)
pr_number, auth = _auth(h)
remote=remote, ctx = _gather_author_duplicate_work_context(
host=host, issue_number,
org=org, h=h,
repo=repo, o=o,
r=r,
auth=auth,
exclude_branch_name=exclude_branch_name,
) )
return pr_work_lease.assess_conflict_fix_push( result = author_duplicate_work_gate.classify_and_assess(
pr_number=pr_number, ctx["issue"],
comments=comments, stage=stage,
branch_head_before=branch_head_before, comments=ctx["comments"],
branch_head_after=branch_head_after, open_prs=ctx["open_prs"],
worktree_path=worktree_path, branch_names=ctx["branch_names"],
push_cwd=push_cwd, allow_stale_takeover=allow_stale_takeover,
is_fast_forward=is_fast_forward,
) )
assessment = result.get("duplicate_work") or {}
return {
"success": True,
"performed": False,
"issue_number": issue_number,
"stage": stage,
"allowed": assessment.get("allowed"),
"block": assessment.get("block"),
"eligibility_class": assessment.get("eligibility_class"),
"outcome": assessment.get("outcome"),
"linked_open_pr": assessment.get("linked_open_pr"),
"matching_branches": assessment.get("matching_branches"),
"claim_status": assessment.get("claim_status"),
"reasons": assessment.get("reasons"),
"safe_next_action": assessment.get("safe_next_action"),
"claim": result.get("claim"),
}
@mcp.tool() @mcp.tool()
-482
View File
@@ -1,482 +0,0 @@
"""Conflict-fix and reviewer PR work leases (#399, #407 reader).
Structured PR/issue comments prove exclusive phases so author conflict-fix
pushes cannot race reviewer validation/approval/merge on the same head.
"""
from __future__ import annotations
import re
from datetime import datetime, timedelta, timezone
from typing import Any
REVIEWER_LEASE_MARKER = "<!-- mcp-review-lease:v1 -->"
CONFLICT_FIX_LEASE_MARKER = "<!-- mcp-conflict-fix-lease:v1 -->"
_FULL_SHA = re.compile(r"^[0-9a-f]{40}$", re.IGNORECASE)
_FIELD_RE = re.compile(
r"^\s*([a-z_]+)\s*:\s*(.+?)\s*$",
re.IGNORECASE | re.MULTILINE,
)
_TERMINAL_REVIEWER_PHASES = frozenset({"done", "released", "blocked"})
_ACTIVE_REVIEWER_PHASES = frozenset({
"claimed",
"validating",
"approved",
"request-changes",
"merging",
})
_TERMINAL_CONFLICT_FIX_PHASES = frozenset({"released", "blocked", "done"})
_ACTIVE_CONFLICT_FIX_PHASES = frozenset({"claimed", "pushing", "pushed"})
DEFAULT_CONFLICT_FIX_TTL_MINUTES = 120
DEFAULT_REVIEWER_LEASE_TTL_MINUTES = 120
def _parse_timestamp(value: str | None) -> datetime | None:
if not value:
return None
text = value.strip()
if text.endswith("Z"):
text = text[:-1] + "+00:00"
try:
parsed = datetime.fromisoformat(text)
except ValueError:
return None
if parsed.tzinfo is None:
return parsed.replace(tzinfo=timezone.utc)
return parsed.astimezone(timezone.utc)
def _normalize_sha(value: str | None) -> str | None:
text = (value or "").strip().lower()
if not text:
return None
return text if _FULL_SHA.match(text) else None
def _parse_pr_ref(value: str | None) -> int | None:
digits = re.sub(r"[^\d]", "", value or "")
return int(digits) if digits.isdigit() else None
def _parse_marker_comment(body: str, marker: str) -> dict[str, str] | None:
text = body or ""
if marker not in text:
return None
fields: dict[str, str] = {}
for match in _FIELD_RE.finditer(text):
fields[match.group(1).strip().lower()] = match.group(2).strip()
return fields or None
def parse_reviewer_lease_comment(body: str) -> dict[str, Any] | None:
fields = _parse_marker_comment(body, REVIEWER_LEASE_MARKER)
if not fields:
return None
return {
"lease_kind": "reviewer",
"pr_number": _parse_pr_ref(fields.get("pr")),
"issue_number": _parse_pr_ref(fields.get("issue")),
"reviewer_identity": fields.get("reviewer_identity"),
"profile": fields.get("profile"),
"session_id": fields.get("session_id"),
"worktree": fields.get("worktree"),
"phase": (fields.get("phase") or "").strip().lower() or None,
"candidate_head": _normalize_sha(fields.get("candidate_head")),
"target_branch": fields.get("target_branch"),
"target_branch_sha": _normalize_sha(fields.get("target_branch_sha")),
"last_activity": fields.get("last_activity"),
"expires_at": fields.get("expires_at"),
"blocker": fields.get("blocker"),
"raw_fields": fields,
}
def parse_conflict_fix_lease_comment(body: str) -> dict[str, Any] | None:
fields = _parse_marker_comment(body, CONFLICT_FIX_LEASE_MARKER)
if not fields:
return None
ff = (fields.get("fast_forward") or "").strip().lower()
reviewer_active = (fields.get("reviewer_active") or "").strip().lower()
return {
"lease_kind": "conflict_fix",
"pr_number": _parse_pr_ref(fields.get("pr")),
"branch": fields.get("branch"),
"worktree": fields.get("worktree"),
"profile": fields.get("profile"),
"session_id": fields.get("session_id"),
"phase": (fields.get("phase") or "").strip().lower() or None,
"head_before": _normalize_sha(fields.get("head_before")),
"head_after": _normalize_sha(fields.get("head_after")),
"expires_at": fields.get("expires_at"),
"reviewer_active": reviewer_active in {"yes", "true", "1"},
"fast_forward": ff in {"yes", "true", "1"},
"raw_fields": fields,
}
def _comment_entries(comments: list[dict], *, pr_number: int | None) -> list[dict]:
entries: list[dict] = []
for comment in comments or []:
body = comment.get("body") or ""
for parser in (parse_reviewer_lease_comment, parse_conflict_fix_lease_comment):
parsed = parser(body)
if not parsed:
continue
if pr_number is not None and parsed.get("pr_number") not in (None, pr_number):
continue
entries.append({
**parsed,
"comment_id": comment.get("id"),
"author": (comment.get("user") or {}).get("login") or comment.get("author"),
"created_at": comment.get("created_at"),
"updated_at": comment.get("updated_at"),
})
break
return entries
def _lease_expired(lease: dict, *, now: datetime) -> bool:
expires_at = _parse_timestamp(lease.get("expires_at"))
return bool(expires_at and expires_at <= now)
def _lease_phase_active(lease: dict, *, active_phases: frozenset[str]) -> bool:
phase = (lease.get("phase") or "").strip().lower()
if phase in _TERMINAL_REVIEWER_PHASES or phase in _TERMINAL_CONFLICT_FIX_PHASES:
return False
return phase in active_phases or bool(phase and phase not in (
_TERMINAL_REVIEWER_PHASES | _TERMINAL_CONFLICT_FIX_PHASES
))
def find_active_reviewer_lease(
comments: list[dict],
*,
pr_number: int,
now: datetime | None = None,
) -> dict[str, Any] | None:
"""Return the newest unexpired reviewer lease for *pr_number*, if any."""
now = now or datetime.now(timezone.utc)
candidates = [
entry for entry in _comment_entries(comments, pr_number=pr_number)
if entry.get("lease_kind") == "reviewer"
]
for lease in reversed(candidates):
if _lease_expired(lease, now=now):
continue
phase = (lease.get("phase") or "").strip().lower()
if phase in _TERMINAL_REVIEWER_PHASES:
continue
if phase in _ACTIVE_REVIEWER_PHASES or phase:
return lease
return None
def find_active_conflict_fix_lease(
comments: list[dict],
*,
pr_number: int,
now: datetime | None = None,
) -> dict[str, Any] | None:
"""Return the newest unexpired conflict-fix lease for *pr_number*, if any."""
now = now or datetime.now(timezone.utc)
candidates = [
entry for entry in _comment_entries(comments, pr_number=pr_number)
if entry.get("lease_kind") == "conflict_fix"
]
for lease in reversed(candidates):
if _lease_expired(lease, now=now):
continue
phase = (lease.get("phase") or "").strip().lower()
if phase in _TERMINAL_CONFLICT_FIX_PHASES:
continue
if phase in _ACTIVE_CONFLICT_FIX_PHASES or phase:
return lease
return None
def format_conflict_fix_lease_body(
*,
pr_number: int,
branch: str,
worktree: str,
profile: str,
head_before: str,
phase: str = "claimed",
session_id: str = "unknown",
expires_at: datetime | None = None,
reviewer_active: bool = False,
) -> str:
expires = expires_at or (
datetime.now(timezone.utc) + timedelta(minutes=DEFAULT_CONFLICT_FIX_TTL_MINUTES)
)
expires_text = expires.astimezone(timezone.utc).replace(microsecond=0).isoformat().replace(
"+00:00", "Z"
)
lines = [
CONFLICT_FIX_LEASE_MARKER,
f"pr: #{pr_number}",
f"branch: {branch}",
f"worktree: {worktree}",
f"profile: {profile}",
f"session_id: {session_id}",
f"phase: {phase}",
f"head_before: {head_before}",
f"expires_at: {expires_text}",
f"reviewer_active: {'yes' if reviewer_active else 'no'}",
]
return "\n".join(lines)
def assess_head_sha_equality(
reviewed_head_sha: str | None,
live_head_sha: str | None,
) -> dict[str, Any]:
"""Fail closed when reviewed and live PR heads differ."""
reviewed = _normalize_sha(reviewed_head_sha)
live = _normalize_sha(live_head_sha)
reasons: list[str] = []
if not reviewed or not live:
reasons.append(
"reviewed/live head SHA missing or not full 40-hex; fail closed"
)
elif reviewed != live:
reasons.append(
"PR head changed after validation; re-pin and re-validate before "
"approval or merge"
)
proven = not reasons
return {
"proven": proven,
"block": not proven,
"reasons": reasons,
"reviewed_head_sha": reviewed,
"live_head_sha": live,
"head_changed": bool(reviewed and live and reviewed != live),
}
def assess_conflict_fix_push(
*,
pr_number: int,
comments: list[dict],
branch_head_before: str | None,
branch_head_after: str | None,
worktree_path: str | None,
push_cwd: str | None,
is_fast_forward: bool | None,
now: datetime | None = None,
) -> dict[str, Any]:
"""Author pre-push gate: block when a reviewer holds an active lease."""
now = now or datetime.now(timezone.utc)
reasons: list[str] = []
reviewer_lease = find_active_reviewer_lease(comments, pr_number=pr_number, now=now)
conflict_lease = find_active_conflict_fix_lease(comments, pr_number=pr_number, now=now)
if reviewer_lease:
reasons.append(
f"active reviewer lease on PR #{pr_number} "
f"(phase={reviewer_lease.get('phase')}); author push blocked"
)
head_before = _normalize_sha(branch_head_before)
head_after = _normalize_sha(branch_head_after)
if not head_before:
reasons.append("branch head before push missing or invalid SHA")
if head_after and head_before and head_before == head_after:
reasons.append("branch head unchanged; no push to perform")
worktree = (worktree_path or "").strip()
cwd = (push_cwd or "").strip()
if not worktree:
reasons.append("worktree path required for conflict-fix push proof")
elif cwd and worktree and not cwd.rstrip("/").endswith(worktree.rstrip("/").split("/")[-1]):
if worktree not in cwd:
reasons.append(
f"push cwd '{cwd}' does not match session worktree '{worktree}'"
)
if is_fast_forward is False:
reasons.append("non-fast-forward push rejected for conflict-fix (fail closed)")
if conflict_lease and conflict_lease.get("phase") == "pushing":
owner = conflict_lease.get("worktree")
if owner and worktree and owner != worktree:
reasons.append(
f"sibling conflict-fix lease active from worktree '{owner}'"
)
push_allowed = not reasons
return {
"push_allowed": push_allowed,
"block": not push_allowed,
"reasons": reasons,
"active_reviewer_lease": reviewer_lease,
"active_conflict_fix_lease": conflict_lease,
"branch_head_before": head_before,
"branch_head_after": head_after,
"reviewer_was_active": bool(reviewer_lease),
"fast_forward": is_fast_forward,
}
def assess_reviewer_mutation_blocked(
*,
pr_number: int,
comments: list[dict],
reviewed_head_sha: str | None,
live_head_sha: str | None,
mutation: str,
now: datetime | None = None,
) -> dict[str, Any]:
"""Reviewer gate: block when conflict-fix lease active or head moved."""
now = now or datetime.now(timezone.utc)
reasons: list[str] = []
conflict_lease = find_active_conflict_fix_lease(comments, pr_number=pr_number, now=now)
if conflict_lease and (conflict_lease.get("phase") or "") in _ACTIVE_CONFLICT_FIX_PHASES:
reasons.append(
f"active conflict-fix lease on PR #{pr_number} "
f"(phase={conflict_lease.get('phase')}); reviewer {mutation} blocked"
)
head_check = assess_head_sha_equality(reviewed_head_sha, live_head_sha)
if head_check["block"]:
reasons.extend(head_check["reasons"])
if not _normalize_sha(reviewed_head_sha):
reasons.append(
f"reviewed head SHA required before reviewer {mutation} (fail closed)"
)
allowed = not reasons
return {
"mutation_allowed": allowed,
"block": not allowed,
"reasons": reasons,
"active_conflict_fix_lease": conflict_lease,
"head_check": head_check,
"reviewed_head_sha": head_check.get("reviewed_head_sha"),
"live_head_sha": head_check.get("live_head_sha"),
"push_during_validation": bool(
conflict_lease and conflict_lease.get("phase") in {"pushing", "pushed"}
),
}
_REVIEWED_HEAD_RE = re.compile(
r"reviewed head sha\s*:\s*([0-9a-f]{40})",
re.IGNORECASE,
)
_LIVE_HEAD_BEFORE_APPROVAL_RE = re.compile(
r"(?:live head sha before approval|final live head sha before approval)\s*:\s*([0-9a-f]{40})",
re.IGNORECASE,
)
_LIVE_HEAD_BEFORE_MERGE_RE = re.compile(
r"(?:live head sha before merge|final live head sha before merge)\s*:\s*([0-9a-f]{40})",
re.IGNORECASE,
)
_PUSH_DURING_VALIDATION_RE = re.compile(
r"push(?:es)? occurred during validation\s*:\s*(yes|no|true|false)",
re.IGNORECASE,
)
_CONFLICT_HEAD_BEFORE_RE = re.compile(
r"branch head before push\s*:\s*([0-9a-f]{40})",
re.IGNORECASE,
)
_CONFLICT_HEAD_AFTER_RE = re.compile(
r"branch head after push\s*:\s*([0-9a-f]{40})",
re.IGNORECASE,
)
_REVIEWER_LEASE_STATUS_RE = re.compile(
r"active reviewer lease status\s*:\s*(.+)$",
re.IGNORECASE | re.MULTILINE,
)
_FAST_FORWARD_RE = re.compile(
r"whether push was fast-forward\s*:\s*(yes|no|true|false)",
re.IGNORECASE,
)
_REVIEWER_ACTIVE_RE = re.compile(
r"whether any reviewer was active\s*:\s*(yes|no|true|false)",
re.IGNORECASE,
)
def assess_reviewer_stale_head_final_report(report_text: str) -> dict[str, Any]:
"""Final-report proof for reviewed vs live head SHAs (#399 AC 6)."""
text = report_text or ""
reasons: list[str] = []
reviewed = _normalize_sha(_REVIEWED_HEAD_RE.search(text).group(1) if _REVIEWED_HEAD_RE.search(text) else None)
live_approval = _normalize_sha(
_LIVE_HEAD_BEFORE_APPROVAL_RE.search(text).group(1)
if _LIVE_HEAD_BEFORE_APPROVAL_RE.search(text)
else None
)
live_merge = _normalize_sha(
_LIVE_HEAD_BEFORE_MERGE_RE.search(text).group(1)
if _LIVE_HEAD_BEFORE_MERGE_RE.search(text)
else None
)
push_during = _PUSH_DURING_VALIDATION_RE.search(text)
if not reviewed:
reasons.append("reviewed head SHA not stated in final report")
if not live_approval:
reasons.append("final live head SHA before approval not stated")
if not live_merge:
reasons.append("final live head SHA before merge not stated")
if not push_during:
reasons.append("whether push occurred during validation not stated")
elif reviewed and live_approval and reviewed != live_approval:
reasons.append("live head before approval differs from reviewed head SHA")
elif reviewed and live_merge and reviewed != live_merge:
reasons.append("live head before merge differs from reviewed head SHA")
proven = not reasons
return {
"proven": proven,
"block": not proven,
"reasons": reasons,
"reviewed_head_sha": reviewed,
"live_head_sha_before_approval": live_approval,
"live_head_sha_before_merge": live_merge,
"push_during_validation": (push_during.group(1).lower() if push_during else None),
}
def assess_conflict_fix_final_report(report_text: str) -> dict[str, Any]:
"""Final-report proof for conflict-fix push sessions (#399 AC 7)."""
text = report_text or ""
reasons: list[str] = []
head_before = _normalize_sha(
_CONFLICT_HEAD_BEFORE_RE.search(text).group(1)
if _CONFLICT_HEAD_BEFORE_RE.search(text)
else None
)
head_after = _normalize_sha(
_CONFLICT_HEAD_AFTER_RE.search(text).group(1)
if _CONFLICT_HEAD_AFTER_RE.search(text)
else None
)
if not head_before:
reasons.append("branch head before push not stated")
if not head_after:
reasons.append("branch head after push not stated")
if not _REVIEWER_LEASE_STATUS_RE.search(text):
reasons.append("active reviewer lease status not stated")
if not _FAST_FORWARD_RE.search(text):
reasons.append("whether push was fast-forward not stated")
if not _REVIEWER_ACTIVE_RE.search(text):
reasons.append("whether any reviewer was active not stated")
proven = not reasons
return {
"proven": proven,
"block": not proven,
"reasons": reasons,
"branch_head_before": head_before,
"branch_head_after": head_after,
}
+19 -1
View File
@@ -3622,18 +3622,33 @@ def assess_work_issue_mode_isolation(report_text: str) -> dict:
} }
def assess_work_issue_duplicate_prevention_report(report_text, **kwargs):
"""#400: work-issue reports must classify duplicate-work prevention."""
from author_duplicate_work_gate import (
assess_work_issue_duplicate_prevention_report as _assess,
)
return _assess(report_text, **kwargs)
def assess_work_issue_final_report(report_text: str) -> dict: def assess_work_issue_final_report(report_text: str) -> dict:
"""#139: composite verifier for work-issue final reports.""" """#139: composite verifier for work-issue final reports."""
checks = { checks = {
"workflow_source": assess_work_issue_workflow_source(report_text), "workflow_source": assess_work_issue_workflow_source(report_text),
"mode_isolation": assess_work_issue_mode_isolation(report_text), "mode_isolation": assess_work_issue_mode_isolation(report_text),
"duplicate_prevention": assess_work_issue_duplicate_prevention_report(
report_text
),
} }
reasons = [] reasons = []
downgraded = False downgraded = False
for name, result in checks.items(): for name, result in checks.items():
verdict = result.get("verdict") verdict = result.get("verdict")
if verdict in ("missing", "incomplete"): if result.get("block"):
downgraded = True
reasons.extend(result.get("reasons") or [])
elif verdict in ("missing", "incomplete"):
downgraded = True downgraded = True
reasons.extend(result.get("reasons") or []) reasons.extend(result.get("reasons") or [])
elif result.get("downgraded") or not result.get("complete", True): elif result.get("downgraded") or not result.get("complete", True):
@@ -3641,6 +3656,9 @@ def assess_work_issue_final_report(report_text: str) -> dict:
reasons.extend( reasons.extend(
f"{name}: {r}" for r in (result.get("reasons") or []) f"{name}: {r}" for r in (result.get("reasons") or [])
) )
elif result.get("proven") is False:
downgraded = True
reasons.extend(result.get("reasons") or [])
grade = "A" if not downgraded else "downgraded" grade = "A" if not downgraded else "downgraded"
return { return {
@@ -732,24 +732,6 @@ The final report must identify:
* whether same-PR merge continuation was allowed * whether same-PR merge continuation was allowed
* whether the run stopped as required * whether the run stopped as required
## 26B. Conflict-fix lease and stale-head protection (#399)
Before validating, approving, or merging a PR:
1. Check for an active conflict-fix lease on the PR; stop if one is active.
2. Pin `expected_head_sha` before validation and pass it to
`gitea_mark_final_review_decision`, `gitea_submit_pr_review`, and
`gitea_merge_pr`.
3. Re-fetch live PR head immediately before approval and merge; refuse when
live head differs from the reviewed SHA.
Final reports must state:
* reviewed head SHA
* final live head SHA before approval
* final live head SHA before merge
* whether any push occurred during validation
## 27. Merge rules ## 27. Merge rules
Before merge, rerun fresh live checks: Before merge, rerun fresh live checks:
@@ -277,6 +277,24 @@ Do not select an issue based only on memory from a previous session.
Before claiming or working on an issue, check whether there is already an open PR, branch, or active claim for that issue. Before claiming or working on an issue, check whether there is already an open PR, branch, or active claim for that issue.
Run `gitea_assess_author_duplicate_work` at these stages and stop when `block` is true:
* `claim` — before `gitea_mark_issue`
* `lock` — before `gitea_lock_issue`
* `worktree` / `edit` — before creating a worktree or editing files
* `commit` — immediately before `git commit`
* `push` — immediately before `git push`
* `create_pr` — immediately before `gitea_create_pr` (also enforced server-side)
`gitea_mark_issue`, `gitea_lock_issue`, and `gitea_create_pr` enforce the same gate server-side and fail closed.
If a concurrent open PR appears after work begins:
* before commit or push — stop and preserve local work without pushing
* after push but before PR creation — produce a reconciliation handoff instead of opening a PR
Final reports must name the duplicate-work outcome (`duplicate PR prevented`, `duplicate branch prevented`, `duplicate commit prevented`, `duplicate push prevented`, `duplicate work not prevented`, or `reconciliation handoff`).
If an open PR already exists for the issue, do not implement duplicate work. If an open PR already exists for the issue, do not implement duplicate work.
Classify the issue as: Classify the issue as:
@@ -577,29 +595,6 @@ After push, report:
If push fails, stop and produce a recovery handoff. If push fails, stop and produce a recovery handoff.
## 20A. Conflict-fix lease and push gate (#399)
When pushing to an existing PR branch to resolve merge conflicts:
1. Call `gitea_acquire_conflict_fix_lease` before any push.
2. Call `gitea_assess_conflict_fix_push` immediately before `git push` with:
* branch head before push
* branch head after push (local)
* session worktree path
* push cwd
* whether the push is fast-forward
3. Do not push when a reviewer holds an active lease on the same PR.
4. Do not force-push.
5. Do not push from the main checkout or wrong cwd.
Conflict-fix final reports must state:
* branch head before push
* branch head after push
* active reviewer lease status
* whether push was fast-forward
* whether any reviewer was active
## 21. PR creation rules ## 21. PR creation rules
Create a PR only if implementation and validation pass, unless project policy explicitly allows draft PRs with documented validation failures. Create a PR only if implementation and validation pass, unless project policy explicitly allows draft PRs with documented validation failures.
+102
View File
@@ -0,0 +1,102 @@
"""Tests for early author duplicate-work gate (#400)."""
import sys
import unittest
from pathlib import Path
sys.path.insert(0, str(Path(__file__).resolve().parent.parent))
from author_duplicate_work_gate import ( # noqa: E402
ELIGIBILITY_OPEN_PR_EXISTS,
assess_author_duplicate_work,
assess_work_issue_duplicate_prevention_report,
)
def _open_pr(number: int = 397, issue: int = 395) -> dict:
return {
"number": number,
"head": {"ref": f"feat/issue-{issue}-example"},
"title": f"feat: example (Closes #{issue})",
"body": f"Closes #{issue}",
}
class TestAuthorDuplicateWorkGate(unittest.TestCase):
def test_clear_when_no_duplicates(self):
result = assess_author_duplicate_work(
400,
stage="claim",
open_prs=[],
branch_names=["master", "feat/issue-399-other"],
)
self.assertTrue(result["allowed"])
self.assertFalse(result["block"])
def test_open_pr_blocks_claim(self):
result = assess_author_duplicate_work(
395,
stage="claim",
open_prs=[_open_pr()],
branch_names=[],
)
self.assertFalse(result["allowed"])
self.assertEqual(result["eligibility_class"], ELIGIBILITY_OPEN_PR_EXISTS)
def test_matching_branch_blocks_lock_not_create_pr(self):
branches = ["feat/issue-400-early-duplicate-work-gate"]
lock = assess_author_duplicate_work(
400,
stage="lock",
open_prs=[],
branch_names=branches,
)
self.assertFalse(lock["allowed"])
create_pr = assess_author_duplicate_work(
400,
stage="create_pr",
open_prs=[],
branch_names=branches,
matching_branches=[],
)
self.assertTrue(create_pr["allowed"])
def test_open_pr_blocks_create_pr_stage(self):
result = assess_author_duplicate_work(
395,
stage="create_pr",
open_prs=[_open_pr()],
branch_names=["feat/issue-395-proof-backed-review-handoff"],
)
self.assertFalse(result["allowed"])
self.assertEqual(result["outcome"], "duplicate_pr_prevented")
def test_push_stage_blocks_on_concurrent_pr(self):
result = assess_author_duplicate_work(
395,
stage="push",
open_prs=[_open_pr()],
branch_names=[],
)
self.assertFalse(result["allowed"])
self.assertEqual(result["outcome"], "duplicate_push_prevented")
def test_duplicate_prevention_report_requires_outcome(self):
bad = assess_work_issue_duplicate_prevention_report(
"Duplicate work detected for issue #395."
)
self.assertFalse(bad["proven"])
good = assess_work_issue_duplicate_prevention_report(
"Duplicate PR prevented; reconciliation handoff produced."
)
self.assertTrue(good["proven"])
def test_exported_from_review_proofs(self):
from review_proofs import assess_work_issue_duplicate_prevention_report as exported
self.assertTrue(callable(exported))
if __name__ == "__main__":
unittest.main()
+6
View File
@@ -95,6 +95,12 @@ def test_create_issue_workflow_contract():
assert "## 9. Duplicate search before mutation" in text assert "## 9. Duplicate search before mutation" in text
def test_author_duplicate_work_gate_exported():
from review_proofs import assess_work_issue_duplicate_prevention_report
assert callable(assess_work_issue_duplicate_prevention_report)
def test_work_issue_workflow_contract(): def test_work_issue_workflow_contract():
text = (SKILL_DIR / "workflows" / "work-issue.md").read_text(encoding="utf-8") text = (SKILL_DIR / "workflows" / "work-issue.md").read_text(encoding="utf-8")
assert "canonical: true" in text assert "canonical: true" in text
-207
View File
@@ -1,207 +0,0 @@
#!/usr/bin/env python3
"""Regression tests for conflict-fix and reviewer PR work leases (#399)."""
from __future__ import annotations
import os
import sys
import unittest
from datetime import datetime, timedelta, timezone
sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.abspath(__file__))))
from pr_work_lease import ( # noqa: E402
CONFLICT_FIX_LEASE_MARKER,
REVIEWER_LEASE_MARKER,
assess_conflict_fix_final_report,
assess_conflict_fix_push,
assess_head_sha_equality,
assess_reviewer_mutation_blocked,
assess_reviewer_stale_head_final_report,
format_conflict_fix_lease_body,
parse_conflict_fix_lease_comment,
parse_reviewer_lease_comment,
)
HEAD_A = "a" * 40
HEAD_B = "b" * 40
NOW = datetime(2026, 7, 7, 15, 0, tzinfo=timezone.utc)
def _reviewer_lease_body(*, phase: str = "validating", expires_minutes: int = 60) -> str:
expires = (NOW + timedelta(minutes=expires_minutes)).isoformat().replace("+00:00", "Z")
return "\n".join([
REVIEWER_LEASE_MARKER,
"pr: #376",
"phase: " + phase,
f"candidate_head: {HEAD_A}",
f"expires_at: {expires}",
"profile: prgs-reviewer",
])
def _conflict_fix_body(*, phase: str = "claimed", worktree: str = "branches/fix-376") -> str:
expires = (NOW + timedelta(minutes=60)).isoformat().replace("+00:00", "Z")
return "\n".join([
CONFLICT_FIX_LEASE_MARKER,
"pr: #376",
f"phase: {phase}",
f"worktree: {worktree}",
f"head_before: {HEAD_A}",
f"expires_at: {expires}",
"profile: prgs-author",
])
class TestLeaseParsing(unittest.TestCase):
def test_parse_reviewer_lease(self):
parsed = parse_reviewer_lease_comment(_reviewer_lease_body())
self.assertEqual(parsed["pr_number"], 376)
self.assertEqual(parsed["phase"], "validating")
self.assertEqual(parsed["candidate_head"], HEAD_A)
def test_parse_conflict_fix_lease(self):
parsed = parse_conflict_fix_lease_comment(_conflict_fix_body())
self.assertEqual(parsed["pr_number"], 376)
self.assertEqual(parsed["phase"], "claimed")
class TestConflictFixPushGate(unittest.TestCase):
def test_blocks_push_during_active_reviewer_lease(self):
comments = [{"body": _reviewer_lease_body()}]
result = assess_conflict_fix_push(
pr_number=376,
comments=comments,
branch_head_before=HEAD_A,
branch_head_after=HEAD_B,
worktree_path="branches/fix-376",
push_cwd="/proj/branches/fix-376",
is_fast_forward=True,
now=NOW,
)
self.assertFalse(result["push_allowed"])
self.assertTrue(any("reviewer lease" in r for r in result["reasons"]))
def test_rejects_non_fast_forward(self):
result = assess_conflict_fix_push(
pr_number=376,
comments=[],
branch_head_before=HEAD_A,
branch_head_after=HEAD_B,
worktree_path="branches/fix-376",
push_cwd="/proj/branches/fix-376",
is_fast_forward=False,
now=NOW,
)
self.assertFalse(result["push_allowed"])
self.assertTrue(any("non-fast-forward" in r for r in result["reasons"]))
def test_wrong_cwd_push_attempt(self):
result = assess_conflict_fix_push(
pr_number=376,
comments=[],
branch_head_before=HEAD_A,
branch_head_after=HEAD_B,
worktree_path="branches/fix-376",
push_cwd="/proj/master",
is_fast_forward=True,
now=NOW,
)
self.assertFalse(result["push_allowed"])
self.assertTrue(any("cwd" in r.lower() for r in result["reasons"]))
def test_sibling_conflict_fix_collision(self):
comments = [{"body": _conflict_fix_body(phase="pushing", worktree="branches/other")}]
result = assess_conflict_fix_push(
pr_number=376,
comments=comments,
branch_head_before=HEAD_A,
branch_head_after=HEAD_B,
worktree_path="branches/fix-376",
push_cwd="/proj/branches/fix-376",
is_fast_forward=True,
now=NOW,
)
self.assertFalse(result["push_allowed"])
self.assertTrue(any("sibling conflict-fix" in r for r in result["reasons"]))
class TestReviewerMutationGate(unittest.TestCase):
def test_blocks_review_during_conflict_fix(self):
comments = [{"body": _conflict_fix_body(phase="pushing")}]
result = assess_reviewer_mutation_blocked(
pr_number=376,
comments=comments,
reviewed_head_sha=HEAD_A,
live_head_sha=HEAD_A,
mutation="approve",
now=NOW,
)
self.assertFalse(result["mutation_allowed"])
self.assertTrue(any("conflict-fix lease" in r for r in result["reasons"]))
def test_stale_head_blocks_approval(self):
result = assess_reviewer_mutation_blocked(
pr_number=376,
comments=[],
reviewed_head_sha=HEAD_A,
live_head_sha=HEAD_B,
mutation="merge",
now=NOW,
)
self.assertFalse(result["mutation_allowed"])
self.assertTrue(result["head_check"]["head_changed"])
def test_head_equality_required_fields(self):
result = assess_head_sha_equality(HEAD_A, HEAD_B)
self.assertFalse(result["proven"])
self.assertTrue(result["head_changed"])
class TestFinalReportProof(unittest.TestCase):
def test_reviewer_stale_head_report_requires_fields(self):
result = assess_reviewer_stale_head_final_report("no head proof here")
self.assertFalse(result["proven"])
def test_reviewer_stale_head_report_passes(self):
report = "\n".join([
f"Reviewed head SHA: {HEAD_A}",
f"Final live head SHA before approval: {HEAD_A}",
f"Final live head SHA before merge: {HEAD_A}",
"Push occurred during validation: no",
])
result = assess_reviewer_stale_head_final_report(report)
self.assertTrue(result["proven"])
def test_conflict_fix_report_requires_fields(self):
result = assess_conflict_fix_final_report("incomplete")
self.assertFalse(result["proven"])
def test_conflict_fix_report_passes(self):
report = "\n".join([
f"Branch head before push: {HEAD_A}",
f"Branch head after push: {HEAD_B}",
"Active reviewer lease status: none",
"Whether push was fast-forward: yes",
"Whether any reviewer was active: no",
])
result = assess_conflict_fix_final_report(report)
self.assertTrue(result["proven"])
class TestFormatLease(unittest.TestCase):
def test_format_conflict_fix_lease_includes_marker(self):
body = format_conflict_fix_lease_body(
pr_number=376,
branch="feat/x",
worktree="branches/fix-376",
profile="prgs-author",
head_before=HEAD_A,
)
self.assertIn(CONFLICT_FIX_LEASE_MARKER, body)
parsed = parse_conflict_fix_lease_comment(body)
self.assertEqual(parsed["pr_number"], 376)
if __name__ == "__main__":
unittest.main()