Replace the Superpowers plan pipeline with grill-with-docs, specs, local tickets, and ticket-native execution. BREAKING CHANGE: Remove the legacy Plan CLI, prompt templates, and Superpowers skills.
1762 lines
66 KiB
Python
1762 lines
66 KiB
Python
#!/usr/bin/env python3
|
|
import argparse
|
|
from contextlib import contextmanager
|
|
from dataclasses import dataclass
|
|
from datetime import datetime, timezone
|
|
import json
|
|
import os
|
|
import re
|
|
import subprocess
|
|
import sys
|
|
import tempfile
|
|
import threading
|
|
import time
|
|
from pathlib import Path
|
|
from typing import Any, Iterator, Optional
|
|
|
|
try:
|
|
import fcntl
|
|
except ImportError: # pragma: no cover
|
|
fcntl = None
|
|
|
|
try:
|
|
import msvcrt
|
|
except ImportError: # pragma: no cover
|
|
msvcrt = None
|
|
|
|
|
|
QUEUE_START = "<!-- main-loop:queue:start -->"
|
|
QUEUE_END = "<!-- main-loop:queue:end -->"
|
|
TICKET_STATE_START = "<!-- main-loop:ticket-state:start -->"
|
|
TICKET_STATE_END = "<!-- main-loop:ticket-state:end -->"
|
|
FEATURE_STATE_FILENAME = ".main-loop.json"
|
|
CLAIM_STALE_AFTER_SECONDS = 30 * 60
|
|
TICKET_FILE_RE = re.compile(r"^(?P<number>\d{2,})-(?P<slug>[a-z0-9][a-z0-9-]*)\.md$")
|
|
TITLE_RE = re.compile(
|
|
r"^#\s+(?P<number>\d{2,})\s+[-\N{EN DASH}\N{EM DASH}]\s+(?P<title>\S.*)$",
|
|
re.MULTILINE,
|
|
)
|
|
STATUS_RE = re.compile(r"^\*\*Status:\*\*\s*(?P<status>\S+)\s*$", re.MULTILINE)
|
|
BLOCKED_BY_RE = re.compile(
|
|
r"^\*\*Blocked by:\*\*\s*(?P<blockers>.+?)\s*$", re.MULTILINE
|
|
)
|
|
ACCEPTANCE_RE = re.compile(r"^-\s+\[[ xX]\]\s+\S", re.MULTILINE)
|
|
NO_BLOCKERS_RE = re.compile(
|
|
r"^none(?:\s*[-\N{EN DASH}\N{EM DASH}:]\s*\S.*)?$", re.IGNORECASE
|
|
)
|
|
BLOCKER_SEPARATOR_RE = re.compile(r";|,(?=\s*\d{2,}\b)")
|
|
BLOCKER_ENTRY_RE = re.compile(
|
|
r"^\s*(?P<number>\d{2,})(?:\s*[-\N{EN DASH}\N{EM DASH}:]\s*\S.*)?\s*$"
|
|
)
|
|
QUEUE_ENTRY_RE = re.compile(r"^-\s+`(?P<slug>[a-z0-9][a-z0-9-]*)`\s*$")
|
|
FEATURE_SLUG_RE = re.compile(r"^[a-z0-9][a-z0-9-]*$")
|
|
ALLOWED_STATUSES = {
|
|
"ready-for-agent",
|
|
"claimed",
|
|
"blocked",
|
|
"resolved",
|
|
"skipped",
|
|
}
|
|
SATISFIED_STATUSES = {"resolved", "skipped"}
|
|
THREAD_LOCKS: dict[str, threading.Lock] = {}
|
|
THREAD_LOCKS_GUARD = threading.Lock()
|
|
|
|
|
|
class StateError(ValueError):
|
|
pass
|
|
|
|
|
|
@dataclass(frozen=True)
|
|
class Ticket:
|
|
number: str
|
|
slug: str
|
|
title: str
|
|
blockers: tuple[str, ...]
|
|
status: str
|
|
path: Path
|
|
metadata: dict[str, Any]
|
|
|
|
|
|
@dataclass(frozen=True)
|
|
class Feature:
|
|
slug: str
|
|
path: Path
|
|
tickets: dict[str, Ticket]
|
|
|
|
@property
|
|
def frontier(self) -> list[Ticket]:
|
|
return [
|
|
ticket
|
|
for ticket in self.tickets.values()
|
|
if ticket.status == "ready-for-agent"
|
|
and all(
|
|
self.tickets[blocker].status in SATISFIED_STATUSES
|
|
for blocker in ticket.blockers
|
|
)
|
|
]
|
|
|
|
@property
|
|
def partial(self) -> bool:
|
|
return any(ticket.status == "skipped" for ticket in self.tickets.values())
|
|
|
|
@property
|
|
def state(self) -> str:
|
|
statuses = {ticket.status for ticket in self.tickets.values()}
|
|
if statuses <= SATISFIED_STATUSES:
|
|
return "ready-to-integrate"
|
|
if "claimed" in statuses or statuses & SATISFIED_STATUSES:
|
|
return "active"
|
|
if not self.frontier:
|
|
return "blocked"
|
|
return "queued"
|
|
|
|
|
|
def _thread_lock(lock_path: Path) -> threading.Lock:
|
|
key = str(lock_path.resolve())
|
|
with THREAD_LOCKS_GUARD:
|
|
lock = THREAD_LOCKS.get(key)
|
|
if lock is None:
|
|
lock = threading.Lock()
|
|
THREAD_LOCKS[key] = lock
|
|
return lock
|
|
|
|
|
|
@contextmanager
|
|
def locked_state(state_root: Path) -> Iterator[None]:
|
|
state_root.mkdir(parents=True, exist_ok=True)
|
|
lock_path = state_root / ".main-loop.lock"
|
|
with _thread_lock(lock_path):
|
|
with lock_path.open("a+b") as lock_file:
|
|
if fcntl is not None:
|
|
fcntl.flock(lock_file.fileno(), fcntl.LOCK_EX)
|
|
elif msvcrt is not None: # pragma: no cover
|
|
while True:
|
|
try:
|
|
lock_file.seek(0)
|
|
msvcrt.locking(lock_file.fileno(), msvcrt.LK_LOCK, 1)
|
|
break
|
|
except OSError:
|
|
time.sleep(0.05)
|
|
try:
|
|
yield
|
|
finally:
|
|
if fcntl is not None:
|
|
fcntl.flock(lock_file.fileno(), fcntl.LOCK_UN)
|
|
elif msvcrt is not None: # pragma: no cover
|
|
lock_file.seek(0)
|
|
msvcrt.locking(lock_file.fileno(), msvcrt.LK_UNLCK, 1)
|
|
|
|
|
|
def atomic_write_text(path: Path, text: str) -> None:
|
|
path.parent.mkdir(parents=True, exist_ok=True)
|
|
fd, raw_temp_path = tempfile.mkstemp(
|
|
dir=path.parent, prefix=f".{path.name}.", suffix=".tmp"
|
|
)
|
|
temp_path = Path(raw_temp_path)
|
|
try:
|
|
with os.fdopen(fd, "w", encoding="utf-8", newline="\n") as handle:
|
|
handle.write(text)
|
|
handle.flush()
|
|
os.fsync(handle.fileno())
|
|
os.replace(temp_path, path)
|
|
finally:
|
|
if temp_path.exists():
|
|
temp_path.unlink()
|
|
|
|
|
|
def parse_blockers(raw: str, ticket_path: Path) -> tuple[str, ...]:
|
|
value = raw.strip()
|
|
if NO_BLOCKERS_RE.fullmatch(value):
|
|
return ()
|
|
blockers: list[str] = []
|
|
for entry in BLOCKER_SEPARATOR_RE.split(value):
|
|
match = BLOCKER_ENTRY_RE.fullmatch(entry)
|
|
if not match:
|
|
raise StateError(f"{ticket_path.name}: malformed Blocked by field")
|
|
blockers.append(match.group("number"))
|
|
if len(blockers) != len(set(blockers)):
|
|
raise StateError(f"{ticket_path.name}: duplicate blocker")
|
|
return tuple(blockers)
|
|
|
|
|
|
def parse_ticket_metadata(text: str, ticket_path: Path) -> dict[str, Any]:
|
|
"""Read the scheduler-owned JSON block, if present.
|
|
|
|
The block is deliberately ordinary Markdown so a ticket remains useful when
|
|
opened without the scheduler. Malformed machine state is rejected instead
|
|
of silently losing ownership information.
|
|
"""
|
|
start_count = text.count(TICKET_STATE_START)
|
|
end_count = text.count(TICKET_STATE_END)
|
|
if start_count == 0 and end_count == 0:
|
|
return {}
|
|
if start_count > 1 or end_count > 1:
|
|
raise StateError(f"{ticket_path.name}: multiple ticket state blocks")
|
|
if start_count != 1 or end_count != 1:
|
|
raise StateError(f"{ticket_path.name}: malformed ticket state markers")
|
|
|
|
marker = re.compile(
|
|
rf"{re.escape(TICKET_STATE_START)}\s*\n(?P<body>.*?)\n"
|
|
rf"{re.escape(TICKET_STATE_END)}",
|
|
re.DOTALL,
|
|
)
|
|
match = marker.search(text)
|
|
if not match:
|
|
raise StateError(f"{ticket_path.name}: malformed ticket state")
|
|
body = match.group("body").strip()
|
|
if not body:
|
|
return {}
|
|
try:
|
|
value = json.loads(body)
|
|
except json.JSONDecodeError as exc:
|
|
raise StateError(f"{ticket_path.name}: malformed ticket state") from exc
|
|
if not isinstance(value, dict):
|
|
raise StateError(f"{ticket_path.name}: ticket state must be an object")
|
|
return value
|
|
|
|
|
|
def render_ticket_state(text: str, metadata: dict[str, Any]) -> str:
|
|
body = json.dumps(metadata, ensure_ascii=True, indent=2, sort_keys=True)
|
|
block = f"{TICKET_STATE_START}\n{body}\n{TICKET_STATE_END}"
|
|
marker = re.compile(
|
|
rf"{re.escape(TICKET_STATE_START)}\s*\n.*?\n"
|
|
rf"{re.escape(TICKET_STATE_END)}",
|
|
re.DOTALL,
|
|
)
|
|
if marker.search(text):
|
|
updated = marker.sub(lambda _match: block, text, count=1)
|
|
else:
|
|
separator = "" if not text or text.endswith("\n\n") else "\n"
|
|
if not text.endswith("\n"):
|
|
separator = "\n\n"
|
|
updated = f"{text}{separator}{block}\n"
|
|
return updated
|
|
|
|
|
|
def update_ticket_state(
|
|
ticket: Ticket,
|
|
*,
|
|
status: Optional[str] = None,
|
|
metadata: Optional[dict[str, Any]] = None,
|
|
) -> None:
|
|
if status is not None and status not in ALLOWED_STATUSES:
|
|
raise StateError(f"invalid status {status}")
|
|
text = ticket.path.read_text(encoding="utf-8")
|
|
if status is not None:
|
|
status_pattern = re.compile(r"(?m)^\*\*Status:\*\*\s*[^\r\n]*$")
|
|
if not status_pattern.search(text):
|
|
raise StateError(f"{ticket.path.name}: missing Status")
|
|
text = status_pattern.sub(f"**Status:** {status}", text, count=1)
|
|
if metadata is not None:
|
|
text = render_ticket_state(text, metadata)
|
|
atomic_write_text(ticket.path, text)
|
|
|
|
|
|
def utc_now() -> str:
|
|
return datetime.now(timezone.utc).replace(microsecond=0).isoformat().replace(
|
|
"+00:00", "Z"
|
|
)
|
|
|
|
|
|
def parse_timestamp(raw: str) -> datetime:
|
|
try:
|
|
value = datetime.fromisoformat(raw.replace("Z", "+00:00"))
|
|
except ValueError as exc:
|
|
raise StateError(f"invalid timestamp: {raw}") from exc
|
|
if value.tzinfo is None:
|
|
raise StateError(f"timestamp must include a timezone: {raw}")
|
|
return value.astimezone(timezone.utc)
|
|
|
|
|
|
def format_timestamp(value: datetime) -> str:
|
|
return value.astimezone(timezone.utc).replace(microsecond=0).isoformat().replace(
|
|
"+00:00", "Z"
|
|
)
|
|
|
|
|
|
def claim_is_stale(ticket: Ticket, now: datetime) -> bool:
|
|
heartbeat = ticket.metadata.get("heartbeat_at")
|
|
if not heartbeat:
|
|
raise StateError(f"{ticket.path.name}: claimed ticket has no heartbeat")
|
|
return (
|
|
now - parse_timestamp(str(heartbeat))
|
|
).total_seconds() > CLAIM_STALE_AFTER_SECONDS
|
|
|
|
|
|
def git_run(repo_root: Path, *args: str) -> subprocess.CompletedProcess[str]:
|
|
return subprocess.run(
|
|
["git", *args],
|
|
cwd=repo_root,
|
|
capture_output=True,
|
|
text=True,
|
|
encoding="utf-8",
|
|
errors="replace",
|
|
)
|
|
|
|
|
|
def git_output(repo_root: Path, *args: str) -> str:
|
|
result = git_run(repo_root, *args)
|
|
if result.returncode != 0:
|
|
detail = (result.stderr or result.stdout).strip()
|
|
raise StateError(f"git {' '.join(args)} failed: {detail}")
|
|
return result.stdout.strip()
|
|
|
|
|
|
def resolve_repo_root(repo_root: Path) -> Path:
|
|
candidate = repo_root.resolve()
|
|
if not candidate.exists():
|
|
raise StateError(f"repo root not found: {repo_root}")
|
|
raw = git_output(candidate, "rev-parse", "--show-toplevel")
|
|
return Path(raw).resolve()
|
|
|
|
|
|
def branch_exists(repo_root: Path, branch: str) -> bool:
|
|
return git_run(repo_root, "show-ref", "--verify", "--quiet", f"refs/heads/{branch}").returncode == 0
|
|
|
|
|
|
def current_branch(repo_root: Path) -> str:
|
|
return git_output(repo_root, "branch", "--show-current")
|
|
|
|
|
|
def worktree_branch_paths(repo_root: Path) -> dict[str, Path]:
|
|
output = git_output(repo_root, "worktree", "list", "--porcelain")
|
|
paths: dict[str, Path] = {}
|
|
current_path: Optional[Path] = None
|
|
for line in output.splitlines():
|
|
if line.startswith("worktree "):
|
|
current_path = Path(line[len("worktree ") :]).resolve()
|
|
elif line.startswith("branch refs/heads/") and current_path is not None:
|
|
branch = line[len("branch refs/heads/") :]
|
|
paths[branch] = current_path
|
|
return paths
|
|
|
|
|
|
def checkout_is_dirty(repo_root: Path) -> Optional[str]:
|
|
output = git_output(repo_root, "status", "--porcelain", "--untracked-files=all")
|
|
for line in output.splitlines():
|
|
if len(line) < 4:
|
|
continue
|
|
raw_path = line[3:].strip().strip('"')
|
|
normalized = raw_path.replace("\\", "/")
|
|
if normalized == ".scratch" or normalized.startswith(".scratch/"):
|
|
continue
|
|
return raw_path
|
|
return None
|
|
|
|
|
|
def ensure_git_repository(repo_root: Path) -> None:
|
|
result = git_run(repo_root, "rev-parse", "--git-dir")
|
|
if result.returncode != 0:
|
|
raise StateError(f"not a Git repository: {repo_root}")
|
|
|
|
|
|
def feature_state_path(feature: Feature) -> Path:
|
|
return feature.path / FEATURE_STATE_FILENAME
|
|
|
|
|
|
def load_feature_state(feature: Feature) -> dict[str, Any]:
|
|
path = feature_state_path(feature)
|
|
if not path.exists():
|
|
return {}
|
|
try:
|
|
value = json.loads(path.read_text(encoding="utf-8"))
|
|
except json.JSONDecodeError as exc:
|
|
raise StateError(f"{path.name}: malformed feature state") from exc
|
|
if not isinstance(value, dict):
|
|
raise StateError(f"{path.name}: feature state must be an object")
|
|
return value
|
|
|
|
|
|
def write_feature_state(feature: Feature, state: dict[str, Any]) -> None:
|
|
atomic_write_text(
|
|
feature_state_path(feature),
|
|
json.dumps(state, ensure_ascii=True, indent=2, sort_keys=True) + "\n",
|
|
)
|
|
|
|
|
|
def parse_ticket(path: Path) -> Ticket:
|
|
file_match = TICKET_FILE_RE.fullmatch(path.name)
|
|
if not file_match:
|
|
raise StateError(f"invalid ticket filename: {path.name}")
|
|
text = path.read_text(encoding="utf-8")
|
|
title_match = TITLE_RE.search(text)
|
|
if not title_match:
|
|
raise StateError(f"{path.name}: invalid title")
|
|
number = file_match.group("number")
|
|
if title_match.group("number") != number:
|
|
raise StateError(f"{path.name}: title number does not match filename")
|
|
status_match = STATUS_RE.search(text)
|
|
if not status_match:
|
|
raise StateError(f"{path.name}: missing Status")
|
|
status = status_match.group("status")
|
|
if status not in ALLOWED_STATUSES:
|
|
raise StateError(f"{path.name}: invalid status {status}")
|
|
blocked_by_match = BLOCKED_BY_RE.search(text)
|
|
if not blocked_by_match:
|
|
raise StateError(f"{path.name}: missing Blocked by")
|
|
if not ACCEPTANCE_RE.search(text):
|
|
raise StateError(f"{path.name}: missing acceptance criterion")
|
|
metadata = parse_ticket_metadata(text, path)
|
|
return Ticket(
|
|
number=number,
|
|
slug=file_match.group("slug"),
|
|
title=title_match.group("title").strip(),
|
|
blockers=parse_blockers(blocked_by_match.group("blockers"), path),
|
|
status=status,
|
|
path=path,
|
|
metadata=metadata,
|
|
)
|
|
|
|
|
|
def validate_acyclic(tickets: dict[str, Ticket]) -> None:
|
|
visiting: set[str] = set()
|
|
visited: set[str] = set()
|
|
|
|
def visit(number: str) -> None:
|
|
if number in visiting:
|
|
raise StateError(f"ticket dependency cycle includes {number}")
|
|
if number in visited:
|
|
return
|
|
visiting.add(number)
|
|
for blocker in tickets[number].blockers:
|
|
visit(blocker)
|
|
visiting.remove(number)
|
|
visited.add(number)
|
|
|
|
for number in tickets:
|
|
visit(number)
|
|
|
|
|
|
def load_feature(state_root: Path, slug: str) -> Feature:
|
|
if not FEATURE_SLUG_RE.fullmatch(slug):
|
|
raise StateError(f"invalid feature slug: {slug}")
|
|
feature_path = state_root / slug
|
|
if not (feature_path / "spec.md").is_file():
|
|
raise StateError(f"{slug}: spec.md not found")
|
|
issues_path = feature_path / "issues"
|
|
if not issues_path.is_dir():
|
|
raise StateError(f"{slug}: issues directory not found")
|
|
ticket_paths = sorted(path for path in issues_path.iterdir() if path.is_file())
|
|
if not ticket_paths:
|
|
raise StateError(f"{slug}: no ticket files found")
|
|
tickets: dict[str, Ticket] = {}
|
|
for path in ticket_paths:
|
|
ticket = parse_ticket(path)
|
|
if ticket.number in tickets:
|
|
raise StateError(f"{slug}: duplicate ticket number {ticket.number}")
|
|
tickets[ticket.number] = ticket
|
|
for ticket in tickets.values():
|
|
for blocker in ticket.blockers:
|
|
if blocker not in tickets:
|
|
raise StateError(f"{ticket.path.name}: unknown blocker {blocker}")
|
|
validate_acyclic(tickets)
|
|
return Feature(slug=slug, path=feature_path, tickets=tickets)
|
|
|
|
|
|
def load_queue(queue_path: Path) -> list[str]:
|
|
if not queue_path.exists():
|
|
return []
|
|
lines = queue_path.read_text(encoding="utf-8").splitlines()
|
|
if lines.count(QUEUE_START) != 1 or lines.count(QUEUE_END) != 1:
|
|
raise StateError("queue.md has an invalid managed block")
|
|
try:
|
|
start = lines.index(QUEUE_START)
|
|
end = lines.index(QUEUE_END, start + 1)
|
|
except ValueError as exc:
|
|
raise StateError("queue.md has an invalid managed block") from exc
|
|
slugs: list[str] = []
|
|
for line in lines[start + 1 : end]:
|
|
if not line.strip():
|
|
continue
|
|
match = QUEUE_ENTRY_RE.fullmatch(line.strip())
|
|
if not match:
|
|
raise StateError(f"queue.md has an invalid entry: {line.strip()}")
|
|
slug = match.group("slug")
|
|
if slug in slugs:
|
|
raise StateError(f"queue.md has a duplicate feature: {slug}")
|
|
slugs.append(slug)
|
|
return slugs
|
|
|
|
|
|
def render_queue(slugs: list[str]) -> str:
|
|
entries = [f"- `{slug}`" for slug in slugs]
|
|
return "\n".join(
|
|
["# Feature Queue", "", QUEUE_START, "", *entries, "", QUEUE_END, ""]
|
|
)
|
|
|
|
|
|
def enqueue_feature(state_root: Path, slug: str) -> str:
|
|
queue_path = state_root / "queue.md"
|
|
with locked_state(state_root):
|
|
load_feature(state_root, slug)
|
|
slugs = load_queue(queue_path)
|
|
if slug in slugs:
|
|
return f"EXISTS={slug}"
|
|
slugs.append(slug)
|
|
atomic_write_text(queue_path, render_queue(slugs))
|
|
return f"ENQUEUED={slug}"
|
|
|
|
|
|
def status_report(
|
|
state_root: Path,
|
|
now: Optional[datetime] = None,
|
|
) -> str:
|
|
observed_at = now or datetime.now(timezone.utc)
|
|
with locked_state(state_root):
|
|
slugs = load_queue(state_root / "queue.md")
|
|
if not slugs:
|
|
return "NO FEATURES"
|
|
output: list[str] = []
|
|
for slug in slugs:
|
|
feature = load_feature(state_root, slug)
|
|
feature_metadata = load_feature_state(feature)
|
|
feature_status = (
|
|
"integrated"
|
|
if feature_metadata.get("integration_commit")
|
|
else (
|
|
"blocked"
|
|
if feature_metadata.get("integration_blocked_reason")
|
|
else feature.state
|
|
)
|
|
)
|
|
output.append(
|
|
f"FEATURE={slug} STATE={feature_status} "
|
|
f"PARTIAL={'yes' if feature.partial else 'no'}"
|
|
)
|
|
frontier = "" if feature_status == "integrated" else ",".join(
|
|
ticket.number for ticket in feature.frontier
|
|
)
|
|
output.append(f"FRONTIER={frontier or '-'}")
|
|
if feature_metadata.get("integration_commit"):
|
|
output.append(
|
|
f"MAIN_INTEGRATION_COMMIT={feature_metadata['integration_commit']}"
|
|
)
|
|
elif feature_metadata.get("integration_blocked_reason"):
|
|
output.append(
|
|
f"FEATURE_BLOCKED={slug} "
|
|
f"REASON={feature_metadata['integration_blocked_reason']}"
|
|
)
|
|
for ticket in feature.tickets.values():
|
|
if ticket.status == "claimed":
|
|
stale = claim_is_stale(ticket, observed_at)
|
|
output.append(
|
|
f"CLAIM={slug}/{ticket.number} "
|
|
f"OWNER={ticket.metadata.get('claimed_by')} "
|
|
f"HEARTBEAT={ticket.metadata.get('heartbeat_at')} "
|
|
f"STALE={'yes' if stale else 'no'} "
|
|
f"ISOLATION={ticket.metadata.get('isolation')} "
|
|
f"WORKSPACE={ticket.metadata.get('workspace')}"
|
|
)
|
|
if ticket.metadata.get("last_error"):
|
|
output.append(
|
|
f"TICKET_ERROR={slug}/{ticket.number} "
|
|
f"REASON={ticket.metadata['last_error']}"
|
|
)
|
|
elif ticket.status == "blocked":
|
|
output.append(
|
|
f"BLOCKED={slug}/{ticket.number} "
|
|
f"REASON={ticket.metadata.get('blocked_reason', 'unspecified')}"
|
|
)
|
|
return "\n".join(output)
|
|
|
|
|
|
def render_claim_context(feature: Feature, ticket: Ticket, metadata: dict[str, Any]) -> str:
|
|
values = [
|
|
("FEATURE", feature.slug),
|
|
("TICKET", ticket.number),
|
|
("CONTROL_ROOT", str(metadata["control_root"])),
|
|
("STATE_ROOT", str(metadata["state_root"])),
|
|
("WORKSPACE", str(metadata["workspace"])),
|
|
("BRANCH", str(metadata["ticket_branch"])),
|
|
("BASE", str(metadata["base_commit"])),
|
|
("ISOLATION", str(metadata["isolation"])),
|
|
]
|
|
return "\n".join(f"{key}={value}" for key, value in values)
|
|
|
|
|
|
def active_claims(features: list[Feature]) -> list[tuple[Feature, Ticket]]:
|
|
claims: list[tuple[Feature, Ticket]] = []
|
|
for feature in features:
|
|
for ticket in feature.tickets.values():
|
|
if ticket.status == "claimed":
|
|
if not ticket.metadata.get("claimed_by"):
|
|
raise StateError(f"{ticket.path.name}: claimed ticket has no owner")
|
|
claims.append((feature, ticket))
|
|
return claims
|
|
|
|
|
|
def checkout_branch(repo_root: Path, branch: str) -> None:
|
|
occupied = worktree_branch_paths(repo_root).get(branch)
|
|
if occupied is not None and os.path.normcase(str(occupied)) != os.path.normcase(
|
|
str(repo_root.resolve())
|
|
):
|
|
raise StateError(f"branch {branch} is checked out at {occupied}")
|
|
result = git_run(repo_root, "checkout", branch)
|
|
if result.returncode != 0:
|
|
detail = (result.stderr or result.stdout).strip()
|
|
raise StateError(f"cannot checkout {branch}: {detail}")
|
|
|
|
|
|
def prepare_in_place_claim(
|
|
repo_root: Path,
|
|
feature: Feature,
|
|
ticket: Ticket,
|
|
main_branch: str,
|
|
prior_metadata: Optional[dict[str, Any]] = None,
|
|
) -> dict[str, Any]:
|
|
ensure_git_repository(repo_root)
|
|
branch = current_branch(repo_root)
|
|
if not branch:
|
|
raise StateError("in-place claim requires an attached HEAD")
|
|
feature_branch = f"feature/{feature.slug}"
|
|
ticket_branch = f"ticket/{feature.slug}/{ticket.number}-{ticket.slug}"
|
|
prior_workspace = prior_metadata.get("workspace") if prior_metadata else None
|
|
restores_released_workspace = bool(
|
|
prior_metadata
|
|
and prior_metadata.get("isolation") == "in-place"
|
|
and prior_workspace
|
|
and Path(str(prior_workspace)).resolve() == repo_root.resolve()
|
|
and prior_metadata.get("ticket_branch") == ticket_branch
|
|
and branch == ticket_branch
|
|
)
|
|
dirty_path = checkout_is_dirty(repo_root)
|
|
if dirty_path and not restores_released_workspace:
|
|
raise StateError(f"in-place checkout is dirty: {dirty_path}")
|
|
if not branch_exists(repo_root, main_branch):
|
|
raise StateError(f"main branch not found: {main_branch}")
|
|
|
|
if not branch_exists(repo_root, feature_branch):
|
|
git_output(repo_root, "branch", feature_branch, main_branch)
|
|
feature_head = git_output(repo_root, "rev-parse", feature_branch)
|
|
|
|
if not branch_exists(repo_root, ticket_branch):
|
|
git_output(repo_root, "branch", ticket_branch, feature_branch)
|
|
if branch != ticket_branch:
|
|
checkout_branch(repo_root, ticket_branch)
|
|
control_root = repo_root.resolve()
|
|
return {
|
|
"isolation": "in-place",
|
|
"control_root": str(control_root),
|
|
"workspace": str(control_root),
|
|
"feature_branch": feature_branch,
|
|
"ticket_branch": ticket_branch,
|
|
"base_commit": feature_head,
|
|
}
|
|
|
|
|
|
def add_worktree(
|
|
repo_root: Path,
|
|
workspace: Path,
|
|
branch: str,
|
|
start_point: str,
|
|
) -> None:
|
|
registered = worktree_branch_paths(repo_root).get(branch)
|
|
if registered is not None:
|
|
if os.path.normcase(str(registered)) == os.path.normcase(str(workspace.resolve())):
|
|
if not workspace.is_dir():
|
|
raise StateError(f"registered worktree is missing: {workspace}")
|
|
return
|
|
raise StateError(f"branch {branch} is checked out at {registered}")
|
|
if workspace.exists():
|
|
if any(workspace.iterdir()):
|
|
raise StateError(f"worktree path is not empty: {workspace}")
|
|
else:
|
|
workspace.parent.mkdir(parents=True, exist_ok=True)
|
|
|
|
if branch_exists(repo_root, branch):
|
|
command = ("worktree", "add", str(workspace), branch)
|
|
else:
|
|
command = ("worktree", "add", "-b", branch, str(workspace), start_point)
|
|
result = git_run(repo_root, *command)
|
|
if result.returncode != 0:
|
|
detail = (result.stderr or result.stdout).strip()
|
|
raise StateError(f"cannot create worktree for {branch}: {detail}")
|
|
|
|
|
|
def prepare_worktree_claim(
|
|
repo_root: Path,
|
|
state_root: Path,
|
|
feature: Feature,
|
|
ticket: Ticket,
|
|
main_branch: str,
|
|
prior_metadata: Optional[dict[str, Any]] = None,
|
|
) -> dict[str, Any]:
|
|
ensure_git_repository(repo_root)
|
|
if not branch_exists(repo_root, main_branch):
|
|
raise StateError(f"main branch not found: {main_branch}")
|
|
feature_branch = f"feature/{feature.slug}"
|
|
ticket_branch = f"ticket/{feature.slug}/{ticket.number}-{ticket.slug}"
|
|
if not branch_exists(repo_root, feature_branch):
|
|
git_output(repo_root, "branch", feature_branch, main_branch)
|
|
|
|
worktree_root = state_root / "worktrees"
|
|
feature_state = load_feature_state(feature)
|
|
integration_workspace = Path(
|
|
str(
|
|
feature_state.get(
|
|
"integration_workspace",
|
|
worktree_root / "_integration" / feature.slug,
|
|
)
|
|
)
|
|
).resolve()
|
|
occupied = worktree_branch_paths(repo_root).get(feature_branch)
|
|
if occupied is not None and os.path.normcase(str(occupied)) == os.path.normcase(
|
|
str(repo_root.resolve())
|
|
):
|
|
dirty_path = checkout_is_dirty(repo_root)
|
|
if dirty_path:
|
|
raise StateError(
|
|
f"cannot free feature branch from dirty control checkout: {dirty_path}"
|
|
)
|
|
checkout_branch(repo_root, main_branch)
|
|
add_worktree(repo_root, integration_workspace, feature_branch, main_branch)
|
|
|
|
prior_workspace = prior_metadata.get("workspace") if prior_metadata else None
|
|
ticket_workspace = Path(
|
|
str(
|
|
prior_workspace
|
|
or worktree_root / feature.slug / f"{ticket.number}-{ticket.slug}"
|
|
)
|
|
).resolve()
|
|
feature_head = git_output(repo_root, "rev-parse", feature_branch)
|
|
add_worktree(
|
|
repo_root,
|
|
ticket_workspace,
|
|
ticket_branch,
|
|
feature_branch,
|
|
)
|
|
|
|
feature_state.update(
|
|
{
|
|
"feature_branch": feature_branch,
|
|
"integration_workspace": str(integration_workspace),
|
|
}
|
|
)
|
|
write_feature_state(feature, feature_state)
|
|
return {
|
|
"isolation": "worktree",
|
|
"control_root": str(repo_root.resolve()),
|
|
"workspace": str(ticket_workspace),
|
|
"feature_branch": feature_branch,
|
|
"ticket_branch": ticket_branch,
|
|
"base_commit": feature_head,
|
|
"integration_workspace": str(integration_workspace),
|
|
}
|
|
|
|
|
|
def resume_claim(
|
|
repo_root: Path,
|
|
feature: Feature,
|
|
ticket: Ticket,
|
|
owner: str,
|
|
) -> str:
|
|
metadata = dict(ticket.metadata)
|
|
isolation = metadata.get("isolation")
|
|
workspace = Path(str(metadata.get("workspace", ""))).resolve()
|
|
if isolation == "in-place":
|
|
if workspace != repo_root.resolve():
|
|
raise StateError(
|
|
f"claim workspace mismatch: expected {workspace}, got {repo_root.resolve()}"
|
|
)
|
|
branch = current_branch(repo_root)
|
|
if not branch:
|
|
raise StateError("in-place claim requires an attached HEAD")
|
|
ticket_branch = str(metadata.get("ticket_branch", ""))
|
|
if branch != ticket_branch:
|
|
dirty_path = checkout_is_dirty(repo_root)
|
|
if dirty_path:
|
|
raise StateError(f"in-place checkout is dirty: {dirty_path}")
|
|
checkout_branch(repo_root, ticket_branch)
|
|
elif isolation == "worktree":
|
|
if not workspace.is_dir():
|
|
raise StateError(f"claim workspace not found: {workspace}")
|
|
ticket_branch = str(metadata.get("ticket_branch", ""))
|
|
if current_branch(workspace) != ticket_branch:
|
|
raise StateError(f"claim workspace is not on {ticket_branch}: {workspace}")
|
|
else:
|
|
raise StateError(f"{ticket.path.name}: invalid claim isolation")
|
|
metadata["heartbeat_at"] = utc_now()
|
|
update_ticket_state(ticket, metadata=metadata)
|
|
return render_claim_context(feature, ticket, metadata)
|
|
|
|
|
|
def choose_ticket_for_claim(features: list[Feature]) -> tuple[Optional[Feature], Optional[Ticket], str]:
|
|
for feature in features:
|
|
state = load_feature_state(feature)
|
|
if state.get("integration_commit"):
|
|
continue
|
|
if state.get("integration_blocked_reason"):
|
|
continue
|
|
if feature.state == "ready-to-integrate":
|
|
return feature, None, "integration-required"
|
|
frontier = sorted(feature.frontier, key=lambda item: (int(item.number), item.slug))
|
|
if frontier:
|
|
return feature, frontier[0], "claim"
|
|
if any(ticket.status == "claimed" for ticket in feature.tickets.values()):
|
|
return feature, None, "busy"
|
|
# A feature with no frontier and no active claim is genuinely blocked;
|
|
# later queued features may be developed, but still integrate in order.
|
|
return None, None, "none"
|
|
|
|
|
|
def claim_ticket(
|
|
state_root: Path,
|
|
repo_root: Path,
|
|
owner: str,
|
|
isolation: str,
|
|
main_branch: str,
|
|
) -> str:
|
|
if not owner.strip():
|
|
raise StateError("owner must not be empty")
|
|
state_root = state_root.resolve()
|
|
repo_root = resolve_repo_root(repo_root)
|
|
with locked_state(state_root):
|
|
slugs = load_queue(state_root / "queue.md")
|
|
if not slugs:
|
|
return "NO FEATURES"
|
|
features = [load_feature(state_root, slug) for slug in slugs]
|
|
claims = active_claims(features)
|
|
|
|
for feature, ticket in claims:
|
|
if ticket.metadata.get("claimed_by") == owner:
|
|
return resume_claim(repo_root, feature, ticket, owner)
|
|
|
|
if claims and isolation == "in-place":
|
|
return "BUSY"
|
|
if any(ticket.metadata.get("isolation") == "in-place" for _, ticket in claims):
|
|
return "BUSY"
|
|
|
|
feature, ticket, disposition = choose_ticket_for_claim(features)
|
|
if disposition == "integration-required" and feature is not None:
|
|
return f"INTEGRATION_REQUIRED={feature.slug}"
|
|
if disposition == "busy":
|
|
return "BUSY"
|
|
if ticket is None or feature is None:
|
|
return "NOOP: no claimable tickets"
|
|
|
|
selected_isolation = isolation
|
|
if isolation == "auto":
|
|
selected_isolation = "in-place" if not claims else "worktree"
|
|
try:
|
|
if selected_isolation == "in-place":
|
|
execution = prepare_in_place_claim(
|
|
repo_root, feature, ticket, main_branch, ticket.metadata
|
|
)
|
|
else:
|
|
execution = prepare_worktree_claim(
|
|
repo_root,
|
|
state_root,
|
|
feature,
|
|
ticket,
|
|
main_branch,
|
|
ticket.metadata,
|
|
)
|
|
except (OSError, StateError) as exc:
|
|
now = utc_now()
|
|
reason = f"claim preparation failed: {exc}"
|
|
metadata = dict(ticket.metadata)
|
|
append_history(
|
|
metadata,
|
|
"claim-blocked",
|
|
now,
|
|
owner=owner,
|
|
reason=reason,
|
|
)
|
|
metadata.update(
|
|
{
|
|
"claimed_by": owner,
|
|
"claimed_at": now,
|
|
"blocked_at": now,
|
|
"blocked_reason": reason,
|
|
"requested_isolation": selected_isolation,
|
|
}
|
|
)
|
|
metadata.pop("heartbeat_at", None)
|
|
update_ticket_state(ticket, status="blocked", metadata=metadata)
|
|
raise StateError(reason) from exc
|
|
now = utc_now()
|
|
metadata = dict(ticket.metadata)
|
|
execution["state_root"] = str(state_root)
|
|
metadata.update(execution)
|
|
metadata.update(
|
|
{
|
|
"claimed_by": owner,
|
|
"claimed_at": now,
|
|
"heartbeat_at": now,
|
|
}
|
|
)
|
|
update_ticket_state(ticket, status="claimed", metadata=metadata)
|
|
return render_claim_context(feature, ticket, metadata)
|
|
|
|
|
|
def evidence_field(
|
|
evidence: str,
|
|
key: str,
|
|
label: str,
|
|
expected: Optional[str] = None,
|
|
) -> str:
|
|
matches = re.findall(
|
|
rf"(?:^|[;,])\s*{re.escape(key)}\s*=\s*([^;,]+)",
|
|
evidence,
|
|
re.IGNORECASE,
|
|
)
|
|
if not matches or not matches[0].strip():
|
|
required = f"{key}={expected}" if expected is not None else key
|
|
raise StateError(f"{label} evidence must include {required}")
|
|
if len(matches) != 1:
|
|
raise StateError(f"{label} evidence must include exactly one {key}")
|
|
value = matches[0].strip()
|
|
if expected is not None and value.lower() != expected.lower():
|
|
raise StateError(f"{label} evidence must include {key}={expected}")
|
|
return value
|
|
|
|
|
|
def require_review_passed(reviewed: str) -> tuple[str, str]:
|
|
for axis in ("standards", "spec"):
|
|
evidence_field(reviewed, axis, "review", "pass")
|
|
return (
|
|
evidence_field(reviewed, "commit", "review"),
|
|
evidence_field(reviewed, "base", "review"),
|
|
)
|
|
|
|
|
|
def require_verification_passed(evidence: str, label: str) -> str:
|
|
if not evidence.strip():
|
|
raise StateError(f"{label} evidence is required")
|
|
evidence_field(evidence, "result", label, "pass")
|
|
return evidence_field(evidence, "commit", label)
|
|
|
|
|
|
def normalize_commit(repo_root: Path, commit: str, label: str) -> str:
|
|
result = git_run(repo_root, "cat-file", "-e", f"{commit}^{{commit}}")
|
|
if result.returncode != 0:
|
|
raise StateError(f"{label} commit not found: {commit}")
|
|
return git_output(repo_root, "rev-parse", f"{commit}^{{commit}}")
|
|
|
|
|
|
def is_ancestor(repo_root: Path, ancestor: str, descendant: str) -> bool:
|
|
return git_run(repo_root, "merge-base", "--is-ancestor", ancestor, descendant).returncode == 0
|
|
|
|
|
|
def find_no_ff_merge(
|
|
repo_root: Path,
|
|
head: str,
|
|
merged_head: str,
|
|
label: str,
|
|
) -> tuple[str, str]:
|
|
merges = git_run(repo_root, "rev-list", "--first-parent", "--merges", head)
|
|
if merges.returncode != 0:
|
|
detail = (merges.stderr or merges.stdout).strip()
|
|
raise StateError(f"cannot inspect {label} integration history: {detail}")
|
|
for merge_commit in merges.stdout.splitlines():
|
|
parents = git_output(
|
|
repo_root,
|
|
"rev-list",
|
|
"--parents",
|
|
"-n",
|
|
"1",
|
|
merge_commit,
|
|
).split()
|
|
if len(parents) == 3 and parents[2] == merged_head:
|
|
return merge_commit, parents[1]
|
|
raise StateError(f"cannot recover {label} integration merge")
|
|
|
|
|
|
def append_history(
|
|
metadata: dict[str, Any], event: str, at: str, **details: str
|
|
) -> None:
|
|
raw_history = metadata.get("history", [])
|
|
history = list(raw_history) if isinstance(raw_history, list) else []
|
|
history.append({"event": event, "at": at, **details})
|
|
metadata["history"] = history
|
|
|
|
|
|
def heartbeat_ticket(
|
|
state_root: Path,
|
|
feature_slug: str,
|
|
ticket_number: str,
|
|
owner: str,
|
|
at: datetime,
|
|
) -> str:
|
|
with locked_state(state_root.resolve()):
|
|
feature = load_feature(state_root.resolve(), feature_slug)
|
|
ticket = feature.tickets.get(ticket_number)
|
|
if ticket is None:
|
|
raise StateError(f"{feature_slug}: ticket {ticket_number} not found")
|
|
if ticket.status != "claimed":
|
|
raise StateError(f"ticket {feature_slug}/{ticket_number} is not claimed")
|
|
if ticket.metadata.get("claimed_by") != owner:
|
|
raise StateError(f"ticket {feature_slug}/{ticket_number} is owned by another session")
|
|
metadata = dict(ticket.metadata)
|
|
timestamp = format_timestamp(at)
|
|
metadata["heartbeat_at"] = timestamp
|
|
update_ticket_state(ticket, metadata=metadata)
|
|
return f"HEARTBEAT={feature_slug}/{ticket_number}\nAT={timestamp}"
|
|
|
|
|
|
def validate_reclaim_workspace(repo_root: Path, ticket: Ticket) -> None:
|
|
metadata = ticket.metadata
|
|
workspace = Path(str(metadata.get("workspace", ""))).resolve()
|
|
if not workspace.is_dir():
|
|
raise StateError(f"claim workspace not found: {workspace}")
|
|
isolation = metadata.get("isolation")
|
|
if isolation == "in-place" and workspace != repo_root.resolve():
|
|
raise StateError(
|
|
f"claim workspace mismatch: expected {workspace}, got {repo_root.resolve()}"
|
|
)
|
|
ticket_branch = str(metadata.get("ticket_branch", ""))
|
|
if current_branch(workspace) != ticket_branch:
|
|
raise StateError(f"claim workspace is not on {ticket_branch}: {workspace}")
|
|
|
|
|
|
def reclaim_ticket(
|
|
state_root: Path,
|
|
repo_root: Path,
|
|
feature_slug: str,
|
|
ticket_number: str,
|
|
owner: str,
|
|
at: datetime,
|
|
) -> str:
|
|
state_root = state_root.resolve()
|
|
repo_root = resolve_repo_root(repo_root)
|
|
with locked_state(state_root):
|
|
feature = load_feature(state_root, feature_slug)
|
|
ticket = feature.tickets.get(ticket_number)
|
|
if ticket is None:
|
|
raise StateError(f"{feature_slug}: ticket {ticket_number} not found")
|
|
if ticket.status != "claimed":
|
|
raise StateError(f"ticket {feature_slug}/{ticket_number} is not claimed")
|
|
previous_owner = str(ticket.metadata.get("claimed_by", ""))
|
|
if not previous_owner:
|
|
raise StateError(f"ticket {feature_slug}/{ticket_number} has no owner")
|
|
if previous_owner == owner:
|
|
raise StateError("reclaim requires a different owner")
|
|
if not claim_is_stale(ticket, at):
|
|
raise StateError(f"ticket {feature_slug}/{ticket_number} is not stale")
|
|
validate_reclaim_workspace(repo_root, ticket)
|
|
|
|
metadata = dict(ticket.metadata)
|
|
timestamp = format_timestamp(at)
|
|
append_history(
|
|
metadata,
|
|
"reclaim",
|
|
timestamp,
|
|
previous_owner=previous_owner,
|
|
owner=owner,
|
|
)
|
|
metadata.update(
|
|
{
|
|
"claimed_by": owner,
|
|
"claimed_at": timestamp,
|
|
"heartbeat_at": timestamp,
|
|
}
|
|
)
|
|
update_ticket_state(ticket, metadata=metadata)
|
|
return render_claim_context(feature, ticket, metadata)
|
|
|
|
|
|
def finish_nonresolved_ticket(
|
|
state_root: Path,
|
|
feature_slug: str,
|
|
ticket_number: str,
|
|
owner: str,
|
|
result: str,
|
|
reason: str,
|
|
) -> str:
|
|
state_root = state_root.resolve()
|
|
with locked_state(state_root):
|
|
feature = load_feature(state_root, feature_slug)
|
|
ticket = feature.tickets.get(ticket_number)
|
|
if ticket is None:
|
|
raise StateError(f"{feature_slug}: ticket {ticket_number} not found")
|
|
allowed_source = ticket.status == "claimed" or (
|
|
result == "released" and ticket.status == "blocked"
|
|
)
|
|
if not allowed_source:
|
|
raise StateError(
|
|
f"ticket {feature_slug}/{ticket_number} cannot become {result} "
|
|
f"from {ticket.status}"
|
|
)
|
|
metadata = dict(ticket.metadata)
|
|
if metadata.get("claimed_by") != owner:
|
|
raise StateError(f"ticket {feature_slug}/{ticket_number} is owned by another session")
|
|
if result in {"blocked", "skipped"} and not reason.strip():
|
|
raise StateError(f"{result} requires a reason")
|
|
|
|
timestamp = utc_now()
|
|
append_history(metadata, result, timestamp, owner=owner, reason=reason.strip())
|
|
if result == "blocked":
|
|
status = "blocked"
|
|
metadata["blocked_reason"] = reason.strip()
|
|
metadata["blocked_at"] = timestamp
|
|
elif result == "released":
|
|
status = "ready-for-agent"
|
|
metadata["last_owner"] = owner
|
|
metadata["released_at"] = timestamp
|
|
metadata.pop("blocked_reason", None)
|
|
metadata.pop("blocked_at", None)
|
|
metadata.pop("claimed_by", None)
|
|
metadata.pop("claimed_at", None)
|
|
metadata.pop("heartbeat_at", None)
|
|
else:
|
|
status = "skipped"
|
|
metadata["last_owner"] = owner
|
|
metadata["skipped_reason"] = reason.strip()
|
|
metadata["skipped_at"] = timestamp
|
|
metadata.pop("claimed_by", None)
|
|
metadata.pop("claimed_at", None)
|
|
metadata.pop("heartbeat_at", None)
|
|
update_ticket_state(ticket, status=status, metadata=metadata)
|
|
return f"{result.upper()}={feature_slug}/{ticket_number}"
|
|
|
|
|
|
def finish_resolved_ticket(
|
|
state_root: Path,
|
|
repo_root: Path,
|
|
feature_slug: str,
|
|
ticket_number: str,
|
|
owner: str,
|
|
implementation_commit: str,
|
|
feature_head: str,
|
|
review_base: str,
|
|
verified: str,
|
|
reviewed: str,
|
|
) -> str:
|
|
if not implementation_commit:
|
|
raise StateError("implementation commit is required")
|
|
if not feature_head:
|
|
raise StateError("feature head is required")
|
|
if not review_base:
|
|
raise StateError("review base is required")
|
|
verification_commit = require_verification_passed(verified, "verification")
|
|
review_commit, reviewed_base = require_review_passed(reviewed)
|
|
state_root = state_root.resolve()
|
|
repo_root = resolve_repo_root(repo_root)
|
|
|
|
with locked_state(state_root):
|
|
feature = load_feature(state_root, feature_slug)
|
|
ticket = feature.tickets.get(ticket_number)
|
|
if ticket is None:
|
|
raise StateError(f"{feature_slug}: ticket {ticket_number} not found")
|
|
metadata = dict(ticket.metadata)
|
|
if ticket.status == "resolved" and metadata.get("integration_commit"):
|
|
return (
|
|
f"RESOLVED={feature_slug}/{ticket_number}\n"
|
|
f"INTEGRATION_COMMIT={metadata['integration_commit']}"
|
|
)
|
|
if ticket.status != "claimed":
|
|
raise StateError(f"ticket {feature_slug}/{ticket_number} is not claimed")
|
|
if metadata.get("claimed_by") != owner:
|
|
raise StateError(f"ticket {feature_slug}/{ticket_number} is owned by another session")
|
|
isolation = metadata.get("isolation")
|
|
expected_workspace = Path(str(metadata.get("workspace", ""))).resolve()
|
|
if isolation == "in-place":
|
|
if expected_workspace != repo_root.resolve():
|
|
raise StateError(
|
|
f"claim workspace mismatch: expected {expected_workspace}, "
|
|
f"got {repo_root.resolve()}"
|
|
)
|
|
ticket_workspace = repo_root
|
|
integration_workspace = repo_root
|
|
elif isolation == "worktree":
|
|
if not expected_workspace.is_dir():
|
|
raise StateError(f"ticket workspace not found: {expected_workspace}")
|
|
ticket_workspace = expected_workspace
|
|
feature_state = load_feature_state(feature)
|
|
integration_value = feature_state.get("integration_workspace")
|
|
if not integration_value:
|
|
raise StateError(f"{feature_slug}: integration workspace is not recorded")
|
|
integration_workspace = Path(str(integration_value)).resolve()
|
|
if not integration_workspace.is_dir():
|
|
raise StateError(
|
|
f"integration workspace not found: {integration_workspace}"
|
|
)
|
|
else:
|
|
raise StateError(f"{ticket.path.name}: invalid claim isolation")
|
|
|
|
dirty_path = checkout_is_dirty(ticket_workspace)
|
|
if dirty_path:
|
|
raise StateError(f"ticket checkout is dirty: {dirty_path}")
|
|
|
|
ticket_branch = str(metadata.get("ticket_branch", ""))
|
|
feature_branch = str(metadata.get("feature_branch", ""))
|
|
if not branch_exists(repo_root, ticket_branch):
|
|
raise StateError(f"ticket branch not found: {ticket_branch}")
|
|
if not branch_exists(repo_root, feature_branch):
|
|
raise StateError(f"feature branch not found: {feature_branch}")
|
|
ticket_tip = git_output(repo_root, "rev-parse", ticket_branch)
|
|
implementation_commit = normalize_commit(
|
|
repo_root, implementation_commit, "implementation"
|
|
)
|
|
feature_head = normalize_commit(repo_root, feature_head, "feature head")
|
|
review_base = normalize_commit(repo_root, review_base, "review base")
|
|
verification_commit = normalize_commit(
|
|
repo_root, verification_commit, "verification"
|
|
)
|
|
review_commit = normalize_commit(repo_root, review_commit, "review")
|
|
reviewed_base = normalize_commit(
|
|
repo_root, reviewed_base, "review base evidence"
|
|
)
|
|
if implementation_commit != ticket_tip:
|
|
raise StateError(
|
|
"implementation commit must equal ticket branch HEAD"
|
|
)
|
|
if verification_commit != ticket_tip:
|
|
raise StateError(
|
|
"verification commit must equal ticket branch HEAD"
|
|
)
|
|
if review_commit != ticket_tip:
|
|
raise StateError("review commit must equal ticket branch HEAD")
|
|
if review_base != feature_head:
|
|
raise StateError(
|
|
"review base must equal the verified feature head"
|
|
)
|
|
if reviewed_base != review_base:
|
|
raise StateError(
|
|
"review evidence base must equal --review-base"
|
|
)
|
|
if current_branch(ticket_workspace) != ticket_branch:
|
|
raise StateError(
|
|
f"ticket workspace is not on its branch: {ticket_workspace}"
|
|
)
|
|
|
|
current_feature_head = git_output(repo_root, "rev-parse", feature_branch)
|
|
already_integrated = is_ancestor(repo_root, ticket_tip, current_feature_head)
|
|
recovered_integration_commit: Optional[str] = None
|
|
if already_integrated:
|
|
recovered_integration_commit, integration_base = find_no_ff_merge(
|
|
repo_root,
|
|
current_feature_head,
|
|
ticket_tip,
|
|
"ticket",
|
|
)
|
|
if feature_head != integration_base:
|
|
raise StateError(
|
|
"verified feature head must equal the ticket integration base"
|
|
)
|
|
else:
|
|
if current_feature_head != feature_head:
|
|
return f"RETRY: feature advanced\nFEATURE_HEAD={current_feature_head}"
|
|
if not is_ancestor(repo_root, feature_head, ticket_tip):
|
|
raise StateError(
|
|
"ticket branch must include the verified feature head"
|
|
)
|
|
|
|
if not already_integrated:
|
|
integration_dirty = checkout_is_dirty(integration_workspace)
|
|
if integration_dirty:
|
|
raise StateError(
|
|
f"integration checkout is dirty: {integration_dirty}"
|
|
)
|
|
if current_branch(integration_workspace) != feature_branch:
|
|
checkout_branch(integration_workspace, feature_branch)
|
|
merge = git_run(
|
|
integration_workspace,
|
|
"merge",
|
|
"--no-ff",
|
|
"-m",
|
|
f"Integrate {feature_slug}/{ticket_number}",
|
|
ticket_branch,
|
|
)
|
|
if merge.returncode != 0:
|
|
detail = (merge.stderr or merge.stdout).strip()
|
|
metadata["last_error"] = f"merge failed: {detail}"
|
|
update_ticket_state(ticket, metadata=metadata)
|
|
raise StateError(f"ticket merge failed: {detail}")
|
|
elif current_branch(integration_workspace) != feature_branch:
|
|
checkout_branch(integration_workspace, feature_branch)
|
|
|
|
integration_commit = (
|
|
recovered_integration_commit
|
|
if recovered_integration_commit is not None
|
|
else git_output(repo_root, "rev-parse", feature_branch)
|
|
)
|
|
metadata.update(
|
|
{
|
|
"implementation_commit": implementation_commit,
|
|
"ticket_head": ticket_tip,
|
|
"integration_commit": integration_commit,
|
|
"verified": verified,
|
|
"reviewed": reviewed,
|
|
"resolved_at": utc_now(),
|
|
}
|
|
)
|
|
metadata.pop("last_error", None)
|
|
update_ticket_state(ticket, status="resolved", metadata=metadata)
|
|
return (
|
|
f"RESOLVED={feature_slug}/{ticket_number}\n"
|
|
f"INTEGRATION_COMMIT={integration_commit}"
|
|
)
|
|
|
|
|
|
def branch_workspace(repo_root: Path, branch: str) -> Optional[Path]:
|
|
return worktree_branch_paths(repo_root).get(branch)
|
|
|
|
|
|
def prepare_main_workspace(repo_root: Path, main_branch: str) -> Path:
|
|
workspace = branch_workspace(repo_root, main_branch)
|
|
if workspace is not None:
|
|
dirty_path = checkout_is_dirty(workspace)
|
|
if dirty_path:
|
|
raise StateError(f"main checkout is dirty: {dirty_path}")
|
|
return workspace
|
|
dirty_path = checkout_is_dirty(repo_root)
|
|
if dirty_path:
|
|
raise StateError(f"control checkout is dirty: {dirty_path}")
|
|
if not current_branch(repo_root):
|
|
raise StateError("control checkout has detached HEAD")
|
|
checkout_branch(repo_root, main_branch)
|
|
return repo_root
|
|
|
|
|
|
def path_is_within(path: Path, parent: Path) -> bool:
|
|
try:
|
|
return os.path.commonpath((str(path.resolve()), str(parent.resolve()))) == str(
|
|
parent.resolve()
|
|
)
|
|
except ValueError:
|
|
return False
|
|
|
|
|
|
def remove_clean_worktree(
|
|
repo_root: Path, workspace: Path, allowed_root: Path
|
|
) -> Optional[str]:
|
|
workspace = workspace.resolve()
|
|
if not path_is_within(workspace, allowed_root):
|
|
return f"preserved unexpected worktree path {workspace}"
|
|
if not workspace.exists():
|
|
return None
|
|
dirty_path = checkout_is_dirty(workspace)
|
|
if dirty_path:
|
|
return f"preserved dirty worktree {workspace}: {dirty_path}"
|
|
result = git_run(repo_root, "worktree", "remove", str(workspace))
|
|
if result.returncode != 0:
|
|
detail = (result.stderr or result.stdout).strip()
|
|
return f"could not remove worktree {workspace}: {detail}"
|
|
return None
|
|
|
|
|
|
def cleanup_integrated_feature_worktrees(
|
|
repo_root: Path, state_root: Path, feature: Feature
|
|
) -> list[str]:
|
|
warnings: list[str] = []
|
|
allowed_root = (state_root / "worktrees").resolve()
|
|
for ticket in feature.tickets.values():
|
|
if ticket.status not in SATISFIED_STATUSES:
|
|
continue
|
|
if ticket.metadata.get("isolation") != "worktree":
|
|
continue
|
|
workspace_value = ticket.metadata.get("workspace")
|
|
if not workspace_value:
|
|
continue
|
|
warning = remove_clean_worktree(
|
|
repo_root, Path(str(workspace_value)), allowed_root
|
|
)
|
|
if warning:
|
|
warnings.append(warning)
|
|
|
|
feature_metadata = load_feature_state(feature)
|
|
integration_value = feature_metadata.get("integration_workspace")
|
|
if integration_value:
|
|
warning = remove_clean_worktree(
|
|
repo_root, Path(str(integration_value)), allowed_root
|
|
)
|
|
if warning:
|
|
warnings.append(warning)
|
|
return warnings
|
|
|
|
|
|
def block_feature_integration(state_root: Path, feature_slug: str, reason: str) -> str:
|
|
if not reason.strip():
|
|
raise StateError("feature block requires a reason")
|
|
state_root = state_root.resolve()
|
|
with locked_state(state_root):
|
|
slugs = load_queue(state_root / "queue.md")
|
|
if feature_slug not in slugs:
|
|
raise StateError(f"feature is not queued: {feature_slug}")
|
|
feature = load_feature(state_root, feature_slug)
|
|
state = load_feature_state(feature)
|
|
if state.get("integration_commit"):
|
|
raise StateError(f"feature is already integrated: {feature_slug}")
|
|
if feature.state != "ready-to-integrate":
|
|
raise StateError(
|
|
f"feature {feature_slug} is not ready to integrate: {feature.state}"
|
|
)
|
|
timestamp = utc_now()
|
|
state["integration_blocked_reason"] = reason.strip()
|
|
state["integration_blocked_at"] = timestamp
|
|
append_history(
|
|
state,
|
|
"feature-blocked",
|
|
timestamp,
|
|
reason=reason.strip(),
|
|
)
|
|
write_feature_state(feature, state)
|
|
return f"FEATURE_BLOCKED={feature_slug}"
|
|
|
|
|
|
def release_feature_integration(state_root: Path, feature_slug: str) -> str:
|
|
state_root = state_root.resolve()
|
|
with locked_state(state_root):
|
|
slugs = load_queue(state_root / "queue.md")
|
|
if feature_slug not in slugs:
|
|
raise StateError(f"feature is not queued: {feature_slug}")
|
|
feature = load_feature(state_root, feature_slug)
|
|
state = load_feature_state(feature)
|
|
if not state.get("integration_blocked_reason"):
|
|
raise StateError(f"feature is not blocked: {feature_slug}")
|
|
timestamp = utc_now()
|
|
append_history(state, "feature-released", timestamp)
|
|
state.pop("integration_blocked_reason", None)
|
|
state.pop("integration_blocked_at", None)
|
|
write_feature_state(feature, state)
|
|
return f"FEATURE_RELEASED={feature_slug}"
|
|
|
|
|
|
def integrate_feature(
|
|
state_root: Path,
|
|
repo_root: Path,
|
|
feature_slug: str,
|
|
feature_head: str,
|
|
verified: str,
|
|
main_verified: str,
|
|
reviewed: str,
|
|
main_branch: str,
|
|
allow_partial: bool,
|
|
) -> str:
|
|
if not feature_head:
|
|
raise StateError("feature head is required")
|
|
feature_verification_commit = require_verification_passed(
|
|
verified, "feature verification"
|
|
)
|
|
main_verification_commit = require_verification_passed(
|
|
main_verified, "main candidate verification"
|
|
)
|
|
review_commit, reviewed_base = require_review_passed(reviewed)
|
|
state_root = state_root.resolve()
|
|
repo_root = resolve_repo_root(repo_root)
|
|
|
|
with locked_state(state_root):
|
|
slugs = load_queue(state_root / "queue.md")
|
|
if feature_slug not in slugs:
|
|
raise StateError(f"feature is not queued: {feature_slug}")
|
|
queued_features = [load_feature(state_root, slug) for slug in slugs]
|
|
requested = next(
|
|
feature for feature in queued_features if feature.slug == feature_slug
|
|
)
|
|
requested_state = load_feature_state(requested)
|
|
if requested_state.get("integration_commit"):
|
|
warnings = cleanup_integrated_feature_worktrees(
|
|
repo_root, state_root, requested
|
|
)
|
|
output = [
|
|
f"INTEGRATED={feature_slug}",
|
|
f"MAIN_INTEGRATION_COMMIT={requested_state['integration_commit']}",
|
|
]
|
|
output.extend(f"WARNING={warning}" for warning in warnings)
|
|
return "\n".join(output)
|
|
if requested_state.get("integration_blocked_reason"):
|
|
raise StateError(
|
|
f"feature {feature_slug} is blocked: "
|
|
f"{requested_state['integration_blocked_reason']}"
|
|
)
|
|
|
|
first_pending: Optional[Feature] = None
|
|
for feature in queued_features:
|
|
if not load_feature_state(feature).get("integration_commit"):
|
|
first_pending = feature
|
|
break
|
|
if first_pending is None:
|
|
return "NOOP: all features integrated"
|
|
if first_pending.slug != feature_slug:
|
|
raise StateError(
|
|
f"cannot integrate {feature_slug}: earlier feature "
|
|
f"{first_pending.slug} is not integrated"
|
|
)
|
|
if requested.state != "ready-to-integrate":
|
|
raise StateError(
|
|
f"feature {feature_slug} is not ready to integrate: {requested.state}"
|
|
)
|
|
if requested.partial and not allow_partial:
|
|
raise StateError(
|
|
f"feature {feature_slug} is partial; pass --allow-partial explicitly"
|
|
)
|
|
|
|
feature_branch = str(
|
|
requested_state.get("feature_branch", f"feature/{feature_slug}")
|
|
)
|
|
if not branch_exists(repo_root, feature_branch):
|
|
raise StateError(f"feature branch not found: {feature_branch}")
|
|
if not branch_exists(repo_root, main_branch):
|
|
raise StateError(f"main branch not found: {main_branch}")
|
|
current_feature_head = git_output(repo_root, "rev-parse", feature_branch)
|
|
feature_head = normalize_commit(repo_root, feature_head, "feature head")
|
|
if current_feature_head != feature_head:
|
|
return f"RETRY: feature advanced\nFEATURE_HEAD={current_feature_head}"
|
|
main_head = git_output(repo_root, "rev-parse", main_branch)
|
|
already_integrated = is_ancestor(
|
|
repo_root, current_feature_head, main_head
|
|
)
|
|
if not already_integrated and not is_ancestor(
|
|
repo_root, main_head, current_feature_head
|
|
):
|
|
return f"RETRY: feature needs main sync\nMAIN_HEAD={main_head}"
|
|
|
|
feature_verification_commit = normalize_commit(
|
|
repo_root,
|
|
feature_verification_commit,
|
|
"feature verification",
|
|
)
|
|
main_verification_commit = normalize_commit(
|
|
repo_root,
|
|
main_verification_commit,
|
|
"main candidate verification",
|
|
)
|
|
review_commit = normalize_commit(repo_root, review_commit, "review")
|
|
reviewed_base = normalize_commit(
|
|
repo_root, reviewed_base, "review base evidence"
|
|
)
|
|
if feature_verification_commit != current_feature_head:
|
|
raise StateError(
|
|
"feature verification commit must equal feature HEAD"
|
|
)
|
|
if main_verification_commit != current_feature_head:
|
|
raise StateError(
|
|
"main candidate verification commit must equal feature HEAD"
|
|
)
|
|
if review_commit != current_feature_head:
|
|
raise StateError("review commit must equal feature HEAD")
|
|
recovered_integration_commit: Optional[str] = None
|
|
if already_integrated:
|
|
if current_feature_head == main_head:
|
|
expected_review_base = main_head
|
|
recovered_integration_commit = main_head
|
|
else:
|
|
(
|
|
recovered_integration_commit,
|
|
expected_review_base,
|
|
) = find_no_ff_merge(
|
|
repo_root,
|
|
main_head,
|
|
current_feature_head,
|
|
"feature",
|
|
)
|
|
if reviewed_base != expected_review_base:
|
|
raise StateError(
|
|
"review base must equal the pre-integration main HEAD"
|
|
)
|
|
elif reviewed_base != main_head:
|
|
raise StateError("review base must equal the latest main HEAD")
|
|
|
|
if already_integrated:
|
|
if recovered_integration_commit is None:
|
|
raise StateError("feature integration commit was not recovered")
|
|
integration_commit = recovered_integration_commit
|
|
else:
|
|
main_workspace = prepare_main_workspace(repo_root, main_branch)
|
|
latest_main = git_output(repo_root, "rev-parse", main_branch)
|
|
if latest_main != main_head:
|
|
return f"RETRY: feature needs main sync\nMAIN_HEAD={latest_main}"
|
|
merge = git_run(
|
|
main_workspace,
|
|
"merge",
|
|
"--no-ff",
|
|
"-m",
|
|
f"Integrate feature {feature_slug}",
|
|
feature_branch,
|
|
)
|
|
if merge.returncode != 0:
|
|
detail = (merge.stderr or merge.stdout).strip()
|
|
requested_state["last_error"] = f"main merge failed: {detail}"
|
|
requested_state["integration_blocked_reason"] = (
|
|
f"main merge failed: {detail}"
|
|
)
|
|
requested_state["integration_blocked_at"] = utc_now()
|
|
write_feature_state(requested, requested_state)
|
|
raise StateError(f"feature merge failed: {detail}")
|
|
|
|
integration_commit = git_output(repo_root, "rev-parse", main_branch)
|
|
requested_state.update(
|
|
{
|
|
"feature_branch": feature_branch,
|
|
"feature_head": current_feature_head,
|
|
"integration_commit": integration_commit,
|
|
"verified": verified,
|
|
"main_verified": main_verified,
|
|
"reviewed": reviewed,
|
|
"partial_authorized": bool(allow_partial),
|
|
"integrated_at": utc_now(),
|
|
}
|
|
)
|
|
requested_state.pop("last_error", None)
|
|
write_feature_state(requested, requested_state)
|
|
warnings = cleanup_integrated_feature_worktrees(
|
|
repo_root, state_root, requested
|
|
)
|
|
output = [
|
|
f"INTEGRATED={feature_slug}",
|
|
f"MAIN_INTEGRATION_COMMIT={integration_commit}",
|
|
]
|
|
output.extend(f"WARNING={warning}" for warning in warnings)
|
|
return "\n".join(output)
|
|
|
|
|
|
def build_parser() -> argparse.ArgumentParser:
|
|
parser = argparse.ArgumentParser(prog="main_loop.py")
|
|
subparsers = parser.add_subparsers(dest="command", required=True)
|
|
|
|
enqueue = subparsers.add_parser("enqueue")
|
|
enqueue.add_argument("--state-root", default=".scratch")
|
|
enqueue.add_argument("--feature", required=True)
|
|
|
|
status = subparsers.add_parser("status")
|
|
status.add_argument("--state-root", default=".scratch")
|
|
|
|
claim = subparsers.add_parser("claim")
|
|
claim.add_argument("--state-root", default=".scratch")
|
|
claim.add_argument("--repo-root", default=".")
|
|
claim.add_argument("--owner", required=True)
|
|
claim.add_argument(
|
|
"--isolation", choices=("auto", "in-place", "worktree"), default="auto"
|
|
)
|
|
claim.add_argument("--main-branch", default="main")
|
|
|
|
finish = subparsers.add_parser("finish")
|
|
finish.add_argument("--state-root", default=".scratch")
|
|
finish.add_argument("--repo-root", default=".")
|
|
finish.add_argument("--feature", required=True)
|
|
finish.add_argument("--ticket", required=True)
|
|
finish.add_argument("--owner", required=True)
|
|
finish.add_argument(
|
|
"--result", choices=("resolved", "blocked", "released", "skipped"), required=True
|
|
)
|
|
finish.add_argument("--implementation-commit", default="")
|
|
finish.add_argument("--feature-head", default="")
|
|
finish.add_argument("--review-base", default="")
|
|
finish.add_argument("--verified", default="")
|
|
finish.add_argument("--reviewed", default="")
|
|
finish.add_argument("--reason", default="")
|
|
|
|
heartbeat = subparsers.add_parser("heartbeat")
|
|
heartbeat.add_argument("--state-root", default=".scratch")
|
|
heartbeat.add_argument("--feature", required=True)
|
|
heartbeat.add_argument("--ticket", required=True)
|
|
heartbeat.add_argument("--owner", required=True)
|
|
|
|
reclaim = subparsers.add_parser("reclaim")
|
|
reclaim.add_argument("--state-root", default=".scratch")
|
|
reclaim.add_argument("--repo-root", default=".")
|
|
reclaim.add_argument("--feature", required=True)
|
|
reclaim.add_argument("--ticket", required=True)
|
|
reclaim.add_argument("--owner", required=True)
|
|
|
|
block_feature = subparsers.add_parser("block-feature")
|
|
block_feature.add_argument("--state-root", default=".scratch")
|
|
block_feature.add_argument("--feature", required=True)
|
|
block_feature.add_argument("--reason", required=True)
|
|
|
|
release_feature = subparsers.add_parser("release-feature")
|
|
release_feature.add_argument("--state-root", default=".scratch")
|
|
release_feature.add_argument("--feature", required=True)
|
|
|
|
integrate = subparsers.add_parser("integrate")
|
|
integrate.add_argument("--state-root", default=".scratch")
|
|
integrate.add_argument("--repo-root", default=".")
|
|
integrate.add_argument("--feature", required=True)
|
|
integrate.add_argument("--feature-head", required=True)
|
|
integrate.add_argument("--verified", required=True)
|
|
integrate.add_argument("--main-verified", required=True)
|
|
integrate.add_argument("--reviewed", required=True)
|
|
integrate.add_argument("--main-branch", default="main")
|
|
integrate.add_argument("--allow-partial", action="store_true")
|
|
return parser
|
|
|
|
|
|
def main(argv: list[str]) -> int:
|
|
args = build_parser().parse_args(argv)
|
|
state_root = Path(args.state_root)
|
|
try:
|
|
if args.command == "enqueue":
|
|
message = enqueue_feature(state_root, args.feature)
|
|
elif args.command == "status":
|
|
message = status_report(state_root)
|
|
elif args.command == "claim":
|
|
message = claim_ticket(
|
|
state_root,
|
|
Path(args.repo_root),
|
|
args.owner,
|
|
args.isolation,
|
|
args.main_branch,
|
|
)
|
|
elif args.command == "heartbeat":
|
|
message = heartbeat_ticket(
|
|
state_root,
|
|
args.feature,
|
|
args.ticket,
|
|
args.owner,
|
|
datetime.now(timezone.utc),
|
|
)
|
|
elif args.command == "reclaim":
|
|
message = reclaim_ticket(
|
|
state_root,
|
|
Path(args.repo_root),
|
|
args.feature,
|
|
args.ticket,
|
|
args.owner,
|
|
datetime.now(timezone.utc),
|
|
)
|
|
elif args.command == "block-feature":
|
|
message = block_feature_integration(
|
|
state_root,
|
|
args.feature,
|
|
args.reason,
|
|
)
|
|
elif args.command == "release-feature":
|
|
message = release_feature_integration(
|
|
state_root,
|
|
args.feature,
|
|
)
|
|
elif args.command == "integrate":
|
|
message = integrate_feature(
|
|
state_root,
|
|
Path(args.repo_root),
|
|
args.feature,
|
|
args.feature_head,
|
|
args.verified,
|
|
args.main_verified,
|
|
args.reviewed,
|
|
args.main_branch,
|
|
args.allow_partial,
|
|
)
|
|
elif args.result == "resolved":
|
|
message = finish_resolved_ticket(
|
|
state_root,
|
|
Path(args.repo_root),
|
|
args.feature,
|
|
args.ticket,
|
|
args.owner,
|
|
args.implementation_commit,
|
|
args.feature_head,
|
|
args.review_base,
|
|
args.verified,
|
|
args.reviewed,
|
|
)
|
|
else:
|
|
message = finish_nonresolved_ticket(
|
|
state_root,
|
|
args.feature,
|
|
args.ticket,
|
|
args.owner,
|
|
args.result,
|
|
args.reason,
|
|
)
|
|
except (OSError, StateError, UnicodeError) as exc:
|
|
print(f"ERROR: {exc}", file=sys.stderr)
|
|
return 2
|
|
print(message)
|
|
return 0
|
|
|
|
|
|
if __name__ == "__main__":
|
|
raise SystemExit(main(sys.argv[1:]))
|