#!/usr/bin/env python3 """ s17: Worktree Isolation — git worktree + task-directory binding + event log. Run: python s17_worktree_isolation/code.py Need: pip install anthropic python-dotenv + .env with ANTHROPIC_API_KEY Changes from s16: - Task dataclass gains worktree field (str | None) - validate_worktree_name: reject path traversal and illegal chars - create_worktree: validate name, git worktree add, optional task binding - bind_task_to_worktree: write worktree field only, keep task pending - remove_worktree: safety check before force, no auto-complete - run_git returns (ok, output), events only on success - Teammate tools: + complete_task, run in worktree cwd when bound - scan_unclaimed_tasks: uses can_start() for dependency checking - Idle teammates wait for messages, then scan and claim ready tasks - consume_lead_inbox: unified inbox consumer - 3 new Lead tools: create_worktree, remove_worktree, keep_worktree ASCII topology: Main repo (/) ├── .worktrees/auth/ (branch: wt/auth) ← Task #1 ├── .worktrees/ui/ (branch: wt/ui) ← Task #2 ├── .tasks/task_xxx.json (worktree: "auth") └── .worktrees/events.jsonl """ import os, subprocess, json, time, random, threading, queue, re from pathlib import Path from datetime import datetime from dataclasses import dataclass, asdict, field try: import readline readline.parse_and_bind('set bind-tty-special-chars off') except ImportError: pass from anthropic import Anthropic from dotenv import load_dotenv load_dotenv(override=True) if os.getenv("ANTHROPIC_BASE_URL"): os.environ.pop("ANTHROPIC_AUTH_TOKEN", None) WORKDIR = Path.cwd() client = Anthropic(base_url=os.getenv("ANTHROPIC_BASE_URL")) MODEL = os.environ["MODEL_ID"] # ── Task System (from s12 + s17 worktree field) ── TASKS_DIR = WORKDIR / ".tasks" TASKS_DIR.mkdir(exist_ok=True) task_lock = threading.RLock() @dataclass class Task: id: str subject: str description: str status: str owner: str | None blockedBy: list[str] worktree: str | None = None # s17: bound worktree name def _task_path(task_id: str) -> Path: return TASKS_DIR / f"{task_id}.json" def create_task(subject: str, description: str = "", blockedBy: list[str] | None = None) -> Task: task = Task( id=f"task_{int(time.time())}_{random.randint(0, 9999):04d}", subject=subject, description=description, status="pending", owner=None, blockedBy=blockedBy or [], ) save_task(task) return task def save_task(task: Task): _task_path(task.id).write_text(json.dumps(asdict(task), indent=2)) def load_task(task_id: str) -> Task: return Task(**json.loads(_task_path(task_id).read_text())) def list_tasks() -> list[Task]: return [Task(**json.loads(p.read_text())) for p in sorted(TASKS_DIR.glob("task_*.json"))] def get_task_json(task_id: str) -> str: task = load_task(task_id) return json.dumps(asdict(task), indent=2) def can_start(task_id: str) -> bool: task = load_task(task_id) for dep_id in task.blockedBy: if not _task_path(dep_id).exists(): return False if load_task(dep_id).status != "completed": return False return True def claim_task(task_id: str, owner: str = "agent") -> str: with task_lock: task = load_task(task_id) if task.status != "pending": return f"Task {task_id} is {task.status}, cannot claim" if task.owner: return f"Task {task_id} already owned by {task.owner}" if not can_start(task_id): deps = [d for d in task.blockedBy if (_task_path(d).exists() and load_task(d).status != "completed")] missing = [d for d in task.blockedBy if not _task_path(d).exists()] parts = [] if deps: parts.append(f"blocked by: {deps}") if missing: parts.append(f"missing deps: {missing}") return "Cannot start: " + ", ".join(parts) task.owner = owner task.status = "in_progress" save_task(task) print(f" \033[36m[claim] {task.subject} → in_progress\033[0m") return f"Claimed {task.id} ({task.subject})" def complete_task(task_id: str) -> str: task = load_task(task_id) if task.status != "in_progress": return f"Task {task_id} is {task.status}, cannot complete" task.status = "completed" save_task(task) unblocked = [t.subject for t in list_tasks() if t.status == "pending" and t.blockedBy and can_start(t.id)] print(f" \033[32m[complete] {task.subject} ✓\033[0m") msg = f"Completed {task.id} ({task.subject})" if unblocked: msg += f"\nUnblocked: {', '.join(unblocked)}" return msg # ── Worktree System (s17 new) ── WORKTREES_DIR = WORKDIR / ".worktrees" WORKTREES_DIR.mkdir(exist_ok=True) VALID_WT_NAME = re.compile(r'^[A-Za-z0-9._-]{1,64}$') def validate_worktree_name(name: str) -> str | None: """Return error message if invalid, None if valid.""" if not name: return "Worktree name cannot be empty" if name == "." or name == "..": return f"'{name}' is not a valid worktree name" if not VALID_WT_NAME.match(name): return (f"Invalid worktree name '{name}': " "only letters, digits, dots, underscores, dashes (1-64 chars)") return None def run_git(args: list[str]) -> tuple[bool, str]: """Run git command. Return (ok, output).""" try: r = subprocess.run(["git"] + args, cwd=WORKDIR, capture_output=True, text=True, timeout=30) out = (r.stdout + r.stderr).strip() out = out[:5000] if out else "(no output)" return r.returncode == 0, out except subprocess.TimeoutExpired: return False, "Error: git timeout" def log_event(event_type: str, worktree_name: str, task_id: str = ""): """Append a lifecycle event to events.jsonl.""" event = {"type": event_type, "worktree": worktree_name, "task_id": task_id, "ts": time.time()} events_file = WORKTREES_DIR / "events.jsonl" with open(events_file, "a") as f: f.write(json.dumps(event) + "\n") def create_worktree(name: str, task_id: str = "") -> str: """Create a git worktree with a dedicated branch. Optionally bind to a task.""" err = validate_worktree_name(name) if err: return f"Error: {err}" path = WORKTREES_DIR / name if path.exists(): return f"Worktree '{name}' already exists at {path}" ok, result = run_git(["worktree", "add", str(path), "-b", f"wt/{name}", "HEAD"]) if not ok: return f"Git error: {result}" if task_id: bind_task_to_worktree(task_id, name) log_event("create", name, task_id) print(f" \033[33m[worktree] created: {name} at {path}\033[0m") return f"Worktree '{name}' created at {path}" def bind_task_to_worktree(task_id: str, worktree_name: str): """Write worktree field to task. Keep status as pending for auto-claim.""" task = load_task(task_id) task.worktree = worktree_name save_task(task) print(f" \033[33m[bind] {task.subject} → worktree:{worktree_name}\033[0m") def _count_worktree_changes(path: Path) -> tuple[int, int]: """Count uncommitted files and commits in a worktree.""" try: r1 = subprocess.run(["git", "status", "--porcelain"], cwd=path, capture_output=True, text=True, timeout=10) files = len([l for l in r1.stdout.strip().splitlines() if l.strip()]) r2 = subprocess.run(["git", "log", "@{push}..HEAD", "--oneline"], cwd=path, capture_output=True, text=True, timeout=10) commits = len([l for l in r2.stdout.strip().splitlines() if l.strip()]) return files, commits except Exception: return -1, -1 def remove_worktree(name: str, discard_changes: bool = False) -> str: """Remove worktree. Refuses if uncommitted changes unless discard_changes.""" err = validate_worktree_name(name) if err: return err path = WORKTREES_DIR / name if not path.exists(): return f"Worktree '{name}' not found" if not discard_changes: files, commits = _count_worktree_changes(path) if files < 0: return (f"Cannot verify worktree '{name}' status. " "Use discard_changes=true to force removal.") if files > 0 or commits > 0: return (f"Worktree '{name}' has {files} uncommitted file(s) " f"and {commits} unpushed commit(s). " "Use discard_changes=true to force removal, " "or keep_worktree to preserve for review.") ok1, _ = run_git(["worktree", "remove", str(path), "--force"]) if not ok1: return f"Failed to remove worktree directory for '{name}'" run_git(["branch", "-D", f"wt/{name}"]) log_event("remove", name) print(f" \033[33m[worktree] removed: {name}\033[0m") return f"Worktree '{name}' removed" def keep_worktree(name: str) -> str: """Keep worktree for manual review. Branch preserved.""" err = validate_worktree_name(name) if err: return err log_event("keep", name) print(f" \033[36m[worktree] kept: {name}\033[0m") return f"Worktree '{name}' kept for review (branch: wt/{name})" # ── Prompt Assembly (from s10) ── PROMPT_SECTIONS = { "identity": "You are a coding agent. Act, don't explain.", "tools": "Available tools: bash, read_file, write_file, " "create_task, list_tasks, get_task, claim_task, complete_task, " "spawn_teammate, send_message, " "request_shutdown, request_plan, review_plan, " "create_worktree, remove_worktree, keep_worktree.", "teams": ( "When parallel work would help, first propose a small team with clear " "responsibilities and wait for the user's confirmation. Do not call " "spawn_teammate before the user confirms." ), "workspace": f"Working directory: {WORKDIR}", "memory": "Relevant memories are injected below when available.", } def assemble_system_prompt(context: dict) -> str: sections = [PROMPT_SECTIONS["identity"], PROMPT_SECTIONS["tools"], PROMPT_SECTIONS["teams"], PROMPT_SECTIONS["workspace"]] if context.get("memories"): sections.append(f"Relevant memories:\n{context['memories']}") return "\n\n".join(sections) _last_context_hash, _last_prompt = None, None def get_system_prompt(context: dict) -> str: global _last_context_hash, _last_prompt h = json.dumps(context, sort_keys=True) if h == _last_context_hash and _last_prompt: return _last_prompt _last_context_hash, _last_prompt = h, assemble_system_prompt(context) return _last_prompt # ── Basic Tools ── def safe_path(p: str, cwd: Path = None) -> Path: base = cwd or WORKDIR path = (base / p).resolve() if not path.is_relative_to(base): raise ValueError(f"Path escapes workspace: {p}") return path def run_bash(command: str, cwd: Path = None) -> str: try: r = subprocess.run(command, shell=True, cwd=cwd or WORKDIR, capture_output=True, text=True, timeout=120) out = (r.stdout + r.stderr).strip() return out[:50000] if out else "(no output)" except subprocess.TimeoutExpired: return "Error: Timeout (120s)" def run_read(path: str, limit: int | None = None, cwd: Path = None) -> str: try: lines = safe_path(path, cwd).read_text().splitlines() if limit and limit < len(lines): lines = lines[:limit] + [f"... ({len(lines) - limit} more lines)"] return "\n".join(lines) except Exception as e: return f"Error: {e}" def run_write(path: str, content: str, cwd: Path = None) -> str: try: fp = safe_path(path, cwd) fp.parent.mkdir(parents=True, exist_ok=True) fp.write_text(content) return f"Wrote {len(content)} bytes to {path}" except Exception as e: return f"Error: {e}" # ── MessageBus (from s15) ── MAILBOX_DIR = WORKDIR / ".mailboxes" MAILBOX_DIR.mkdir(exist_ok=True) MAILBOX_ROOT = MAILBOX_DIR.resolve() VALID_AGENT_NAME = re.compile(r"^[A-Za-z0-9_-]{1,64}$") def is_valid_agent_name(name: str) -> bool: return bool(VALID_AGENT_NAME.fullmatch(name)) class MessageBus: def __init__(self): self._lock = threading.RLock() self._changed = threading.Condition(self._lock) def _path(self, agent: str) -> Path: if not is_valid_agent_name(agent): raise ValueError(f"Invalid mailbox recipient: {agent!r}") path = (MAILBOX_DIR / f"{agent}.jsonl").resolve() if not path.is_relative_to(MAILBOX_ROOT): raise ValueError(f"Mailbox path escapes directory: {agent!r}") return path def _read_unlocked(self, agent: str) -> list[dict]: inbox = self._path(agent) if not inbox.exists(): return [] msgs = [json.loads(line) for line in inbox.read_text().splitlines() if line.strip()] inbox.unlink() return msgs def send(self, from_agent: str, to_agent: str, content: str, msg_type: str = "message", metadata: dict | None = None): msg = {"from": from_agent, "to": to_agent, "content": content, "type": msg_type, "ts": time.time(), "metadata": metadata or {}} with self._changed: with open(self._path(to_agent), "a") as f: f.write(json.dumps(msg, ensure_ascii=False) + "\n") self._changed.notify_all() print(f" \033[33m[bus] {from_agent} → {to_agent}: " f"({msg_type}) {content[:50]}\033[0m") def read_inbox(self, agent: str) -> list[dict]: with self._lock: return self._read_unlocked(agent) def peek(self, agent: str) -> bool: with self._lock: inbox = self._path(agent) return inbox.exists() and inbox.stat().st_size > 0 def wait_for_messages(self, agent: str, timeout: float | None = None) -> list[dict]: deadline = None if timeout is None else time.monotonic() + timeout with self._changed: while not self.peek(agent): remaining = (None if deadline is None else deadline - time.monotonic()) if remaining is not None and remaining <= 0: return [] self._changed.wait(remaining) return self._read_unlocked(agent) BUS = MessageBus() active_teammates: dict[str, str] = {} plan_gates: dict[str, str] = {} plan_request_ids: dict[str, str] = {} team_lock = threading.RLock() # ── Protocol State (from s15) ── @dataclass class ProtocolState: request_id: str type: str sender: str target: str status: str payload: str created_at: float = field(default_factory=time.time) pending_requests: dict[str, ProtocolState] = {} def new_request_id() -> str: while True: request_id = f"req_{random.randint(0, 999999):06d}" if request_id not in pending_requests: return request_id def match_response(response_type: str, request_id: str, approve: bool, from_agent: str, to_agent: str) -> bool: with team_lock: state = pending_requests.get(request_id) if not state: print(f" \033[31m[protocol] unknown request_id: {request_id}\033[0m") return False expected = { "shutdown": "shutdown_response", "plan_approval": "plan_approval_response", }[state.type] if response_type != expected: print(f" \033[31m[protocol] expected {expected}, " f"got {response_type}\033[0m") return False if from_agent != state.target or to_agent != state.sender: print(f" \033[31m[protocol] {request_id} responder mismatch\033[0m") return False if state.status != "pending": return False state.status = "approved" if approve else "rejected" icon = "✓" if approve else "✗" color = "32" if approve else "31" print(f" \033[{color}m[protocol] {state.type} {icon} " f"({request_id}: {state.status})\033[0m") return True def consume_lead_inbox(route_protocol=True) -> list[dict]: msgs = BUS.read_inbox("lead") if route_protocol: for msg in msgs: meta = msg.get("metadata", {}) req_id = meta.get("request_id", "") msg_type = msg.get("type", "") if req_id and msg_type.endswith("_response"): match_response(msg_type, req_id, meta.get("approve", False), msg.get("from", ""), msg.get("to", "")) return msgs def format_team_events(msgs: list[dict]) -> str: lines = [] for msg in msgs: request_id = msg.get("metadata", {}).get("request_id") suffix = f" request_id={request_id}" if request_id else "" lines.append( f"[{msg['type']}{suffix}] {msg['from']}: {msg['content']}" ) return "[Team events]\n" + "\n".join(lines) # ── Autonomous Agent (from s16, + worktree cwd) ── IDLE_SCAN_INTERVAL = 2.0 def scan_unclaimed_tasks() -> list[Task]: return [ task for task in list_tasks() if (task.status == "pending" and task.owner is None and can_start(task.id)) ] def claim_next_task(name: str) -> Task | None: for task in scan_unclaimed_tasks(): result = claim_task(task.id, owner=name) if result.startswith("Claimed "): return load_task(task.id) return None def _last_assistant_text(content) -> str: for block in content: if getattr(block, "type", None) == "text": return block.text.strip() if isinstance(block, dict) and block.get("type") == "text": return str(block.get("text", "")).strip() return "" def _run_teammate_tool(name: str, block, handlers: dict) -> str: gate = plan_gates.get(name, "not_required") if (block.name in {"bash", "write_file"} and gate not in {"not_required", "approved"}): return f"Blocked: plan status is {gate}." handler = handlers.get(block.name) return str(handler(**block.input)) if handler else f"Unknown tool: {block.name}" def apply_plan_response(name: str, msg: dict) -> tuple[bool, str]: """Apply only the Lead response for this teammate's current plan.""" metadata = msg.get("metadata", {}) request_id = metadata.get("request_id", "") with team_lock: state = pending_requests.get(request_id) expected_id = plan_request_ids.get(name) valid = ( msg.get("from") == "lead" and msg.get("to") == name and request_id == expected_id and state is not None and state.type == "plan_approval" and state.sender == name and state.target == "lead" and state.status in {"approved", "rejected"} and metadata.get("approve", False) == (state.status == "approved") ) if not valid: return False, "[Ignored plan response: request mismatch]" plan_gates[name] = state.status active_teammates[name] = "working" plan_request_ids.pop(name, None) outcome = state.status return True, f"[Plan {outcome}] {msg['content']}" def apply_shutdown_request(name: str, msg: dict) -> tuple[bool, str]: """Accept only a pending shutdown request sent by Lead to this teammate.""" request_id = msg.get("metadata", {}).get("request_id", "") with team_lock: state = pending_requests.get(request_id) valid = ( msg.get("from") == "lead" and msg.get("to") == name and state is not None and state.type == "shutdown" and state.sender == "lead" and state.target == name and state.status == "pending" and active_teammates.get(name) != "stopping" ) if not valid: return False, "[Ignored shutdown request: request mismatch]" active_teammates[name] = "stopping" return True, request_id def _teammate_send_message(from_name: str, to: str, content: str) -> str: with team_lock: if to != "lead" and to not in active_teammates: return f"Agent '{to}' is not active" BUS.send(from_name, to, content) return f"Sent to {to}" # ── Teammate Thread ── def spawn_teammate_thread(name: str, role: str, prompt: str) -> str: if not is_valid_agent_name(name): return ("Invalid teammate name: use 1-64 letters, digits, " "underscores, or dashes") with team_lock: if name in active_teammates: return f"Teammate '{name}' already exists" active_teammates[name] = "working" plan_gates[name] = "not_required" system = (f"You are '{name}', a {role}. " "Use tools to complete tasks. " "You can list and claim tasks from the board. " "If a task has a worktree, work in that directory. " "When asked for a plan, submit it before bash or write_file " "and wait for approval.") def handle_inbox_message(name: str, msg: dict, messages: list): msg_type = msg.get("type", "message") meta = msg.get("metadata", {}) req_id = meta.get("request_id", "") if msg_type == "shutdown_request": accepted, notice = apply_shutdown_request(name, msg) if not accepted: messages.append({"role": "user", "content": notice}) return False req_id = notice BUS.send(name, "lead", "Shutting down gracefully.", "shutdown_response", {"request_id": req_id, "approve": True}) print(f" \033[35m[protocol] {name} approved shutdown " f"({req_id})\033[0m") return True if msg_type == "plan_approval_response": _, notice = apply_plan_response(name, msg) messages.append({"role": "user", "content": notice}) elif msg_type == "plan_request": messages.append({"role": "user", "content": f"[Plan required] {msg['content']}"}) elif msg_type == "message": messages.append({"role": "user", "content": f"[Message from {msg['from']}] {msg['content']}"}) return False def run(): # Track current worktree for this teammate's cwd wt_ctx = {"path": None} def _wt_cwd() -> Path | None: p = wt_ctx["path"] return Path(p) if p else None def _run_bash(command: str) -> str: return run_bash(command, cwd=_wt_cwd()) def _run_read(path: str) -> str: return run_read(path, cwd=_wt_cwd()) def _run_write(path: str, content: str) -> str: return run_write(path, content, cwd=_wt_cwd()) def _run_list_tasks(): tasks = list_tasks() if not tasks: return "No tasks." return "\n".join( f" {t.id}: {t.subject} [{t.status}]" + (f" (wt:{t.worktree})" if t.worktree else "") for t in tasks) def _run_claim_task(task_id: str): result = claim_task(task_id, owner=name) if "Claimed" in result: # Set worktree cwd if task has one task = load_task(task_id) if task.worktree: wt_ctx["path"] = str(WORKTREES_DIR / task.worktree) else: wt_ctx["path"] = None return result def _run_complete_task(task_id: str): result = complete_task(task_id) wt_ctx["path"] = None return result messages = [{"role": "user", "content": prompt}] sub_tools = [ {"name": "bash", "description": "Run a shell command.", "input_schema": {"type": "object", "properties": {"command": {"type": "string"}}, "required": ["command"]}}, {"name": "read_file", "description": "Read file.", "input_schema": {"type": "object", "properties": {"path": {"type": "string"}}, "required": ["path"]}}, {"name": "write_file", "description": "Write file.", "input_schema": {"type": "object", "properties": {"path": {"type": "string"}, "content": {"type": "string"}}, "required": ["path", "content"]}}, {"name": "send_message", "description": "Send message to another agent.", "input_schema": {"type": "object", "properties": {"to": {"type": "string"}, "content": {"type": "string"}}, "required": ["to", "content"]}}, {"name": "submit_plan", "description": "Submit a plan for Lead approval.", "input_schema": {"type": "object", "properties": {"plan": {"type": "string"}}, "required": ["plan"]}}, {"name": "list_tasks", "description": "List all tasks on the board.", "input_schema": {"type": "object", "properties": {}, "required": []}}, {"name": "claim_task", "description": "Claim a pending task.", "input_schema": {"type": "object", "properties": {"task_id": {"type": "string"}}, "required": ["task_id"]}}, {"name": "complete_task", "description": "Mark an in-progress task as completed.", "input_schema": {"type": "object", "properties": {"task_id": {"type": "string"}}, "required": ["task_id"]}}, ] sub_handlers = { "bash": _run_bash, "read_file": _run_read, "write_file": _run_write, "send_message": lambda to, content: _teammate_send_message( name, to, content), "submit_plan": lambda plan: _teammate_submit_plan(name, plan), "list_tasks": _run_list_tasks, "claim_task": _run_claim_task, "complete_task": _run_complete_task, } # Outer loop: WORK → IDLE cycle while True: if len(messages) <= 3: messages.insert(0, {"role": "user", "content": f"You are '{name}', role: {role}. " f"Continue your work."}) # WORK phase should_shutdown = False for _ in range(10): inbox = BUS.read_inbox(name) for msg in inbox: stopped = handle_inbox_message(name, msg, messages) if stopped: should_shutdown = True break if should_shutdown: break try: response = client.messages.create( model=MODEL, system=system, messages=messages[-20:], tools=sub_tools, max_tokens=8000) except Exception: break messages.append({"role": "assistant", "content": response.content}) if response.stop_reason != "tool_use": summary = _last_assistant_text(response.content) gate = plan_gates.get(name, "not_required") if gate != "pending" and summary: BUS.send(name, "lead", summary, "result") if gate == "pending": with team_lock: active_teammates[name] = "waiting_approval" else: with team_lock: active_teammates[name] = "idle" BUS.send(name, "lead", "Waiting for more work.", "idle_notification") break results = [] for block in response.content: if block.type == "tool_use": output = _run_teammate_tool(name, block, sub_handlers) results.append({"type": "tool_result", "tool_use_id": block.id, "content": str(output)}) messages.append({"role": "user", "content": results}) if should_shutdown: break # IDLE phase: messages take priority, then scan the task board. while True: inbox = BUS.wait_for_messages(name, IDLE_SCAN_INTERVAL) if inbox: for msg in inbox: if handle_inbox_message(name, msg, messages): should_shutdown = True break if should_shutdown or messages[-1]["role"] == "user": break continue task = claim_next_task(name) if not task: continue wt_ctx["path"] = (str(WORKTREES_DIR / task.worktree) if task.worktree else None) workdir = (f"\nWork directory: {wt_ctx['path']}" if wt_ctx["path"] else "") messages.append({ "role": "user", "content": ( f"[Auto-claimed task {task.id}] " f"{task.subject}\n{task.description}{workdir}" ), }) print(f" \033[32m[idle] {name} claimed " f"{task.id}: {task.subject}\033[0m") break if should_shutdown: break with team_lock: active_teammates.pop(name, None) plan_gates.pop(name, None) plan_request_ids.pop(name, None) print(f" \033[32m[teammate] {name} finished\033[0m") threading.Thread(target=run, daemon=True).start() print(f" \033[36m[teammate] {name} spawned as {role}\033[0m") return f"Teammate '{name}' spawned as {role} (autonomous)" def _teammate_submit_plan(from_name: str, plan: str) -> str: with team_lock: if plan_gates.get(from_name) == "pending": return "A plan is already waiting for review." req_id = new_request_id() pending_requests[req_id] = ProtocolState( request_id=req_id, type="plan_approval", sender=from_name, target="lead", status="pending", payload=plan) plan_gates[from_name] = "pending" plan_request_ids[from_name] = req_id active_teammates[from_name] = "waiting_approval" BUS.send(from_name, "lead", plan, "plan_approval_request", {"request_id": req_id}) return f"Plan submitted ({req_id}). Waiting for approval..." # ── Lead Protocol Tools (from s15) ── def run_request_shutdown(teammate: str) -> str: if teammate not in active_teammates: return f"Teammate '{teammate}' is not active" with team_lock: req_id = new_request_id() pending_requests[req_id] = ProtocolState( request_id=req_id, type="shutdown", sender="lead", target=teammate, status="pending", payload="") BUS.send("lead", teammate, "Please shut down gracefully.", "shutdown_request", {"request_id": req_id}) print(f" \033[35m[protocol] shutdown_request → {teammate} " f"({req_id})\033[0m") return f"Shutdown request sent to {teammate} (req: {req_id})" def run_request_plan(teammate: str, task: str) -> str: if teammate not in active_teammates: return f"Teammate '{teammate}' is not active" with team_lock: plan_gates[teammate] = "required" BUS.send("lead", teammate, task, "plan_request") return f"Asked {teammate} to submit a plan" def run_review_plan(request_id: str, approve: bool, feedback: str = "") -> str: with team_lock: state = pending_requests.get(request_id) if not state: return f"Request {request_id} not found" if state.type != "plan_approval": return f"Request {request_id} is not a plan" if state.status != "pending": return f"Request {request_id} already {state.status}" if plan_request_ids.get(state.sender) != request_id: return f"Request {request_id} is not the current plan" state.status = "approved" if approve else "rejected" BUS.send("lead", state.sender, feedback or ("Approved" if approve else "Rejected"), "plan_approval_response", {"request_id": request_id, "approve": approve}) icon = "✓" if approve else "✗" print(f" \033[32m[protocol] plan {icon} ({request_id})\033[0m") return f"Plan {'approved' if approve else 'rejected'} ({request_id})" # ── Lead Worktree Tools (s17 new) ── def run_create_worktree(name: str, task_id: str = "") -> str: return create_worktree(name, task_id) def run_remove_worktree(name: str, discard_changes: bool = False) -> str: return remove_worktree(name, discard_changes) def run_keep_worktree(name: str) -> str: return keep_worktree(name) # ── Basic tool handlers ── def run_create_task(subject: str, description: str = "", blockedBy: list[str] | None = None) -> str: task = create_task(subject, description, blockedBy) deps = f" (blockedBy: {', '.join(blockedBy)})" if blockedBy else "" print(f" \033[34m[create] {task.subject}{deps}\033[0m") return f"Created {task.id}: {task.subject}{deps}" def run_list_tasks() -> str: tasks = list_tasks() if not tasks: return "No tasks." return "\n".join( f" {t.id}: {t.subject} [{t.status}]" + (f" (wt:{t.worktree})" if t.worktree else "") for t in tasks) def run_get_task(task_id: str) -> str: return get_task_json(task_id) def run_claim_task(task_id: str) -> str: return claim_task(task_id, owner="agent") def run_complete_task(task_id: str) -> str: return complete_task(task_id) def run_spawn_teammate(name: str, role: str, prompt: str) -> str: return spawn_teammate_thread(name, role, prompt) def run_send_message(to: str, content: str) -> str: if to not in active_teammates: return f"Teammate '{to}' is not active" BUS.send("lead", to, content) return f"Sent to {to}" # ── Tool Definitions ── TOOLS = [ {"name": "bash", "description": "Run a shell command.", "input_schema": {"type": "object", "properties": {"command": {"type": "string"}}, "required": ["command"]}}, {"name": "read_file", "description": "Read file contents.", "input_schema": {"type": "object", "properties": {"path": {"type": "string"}, "limit": {"type": "integer"}}, "required": ["path"]}}, {"name": "write_file", "description": "Write content to a file.", "input_schema": {"type": "object", "properties": {"path": {"type": "string"}, "content": {"type": "string"}}, "required": ["path", "content"]}}, {"name": "create_task", "description": "Create a task.", "input_schema": {"type": "object", "properties": {"subject": {"type": "string"}, "description": {"type": "string"}, "blockedBy": {"type": "array", "items": {"type": "string"}}}, "required": ["subject"]}}, {"name": "list_tasks", "description": "List all tasks.", "input_schema": {"type": "object", "properties": {}, "required": []}}, {"name": "get_task", "description": "Get full details of a specific task.", "input_schema": {"type": "object", "properties": {"task_id": {"type": "string"}}, "required": ["task_id"]}}, {"name": "claim_task", "description": "Claim a pending task.", "input_schema": {"type": "object", "properties": {"task_id": {"type": "string"}}, "required": ["task_id"]}}, {"name": "complete_task", "description": "Complete an in-progress task.", "input_schema": {"type": "object", "properties": {"task_id": {"type": "string"}}, "required": ["task_id"]}}, {"name": "spawn_teammate", "description": "Spawn an autonomous teammate agent.", "input_schema": {"type": "object", "properties": {"name": { "type": "string", "pattern": "^[A-Za-z0-9_-]{1,64}$", }, "role": {"type": "string"}, "prompt": {"type": "string"}}, "required": ["name", "role", "prompt"]}}, {"name": "send_message", "description": "Send message to a teammate.", "input_schema": {"type": "object", "properties": {"to": {"type": "string"}, "content": {"type": "string"}}, "required": ["to", "content"]}}, {"name": "request_shutdown", "description": "Request a teammate to shut down gracefully.", "input_schema": {"type": "object", "properties": {"teammate": {"type": "string"}}, "required": ["teammate"]}}, {"name": "request_plan", "description": "Ask a teammate to submit a plan for review.", "input_schema": {"type": "object", "properties": {"teammate": {"type": "string"}, "task": {"type": "string"}}, "required": ["teammate", "task"]}}, {"name": "review_plan", "description": "Approve or reject a submitted plan.", "input_schema": {"type": "object", "properties": { "request_id": {"type": "string"}, "approve": {"type": "boolean"}, "feedback": {"type": "string"}}, "required": ["request_id", "approve"]}}, # s17 new: worktree tools {"name": "create_worktree", "description": "Create an isolated git worktree with its own branch.", "input_schema": {"type": "object", "properties": {"name": {"type": "string"}, "task_id": {"type": "string"}}, "required": ["name"]}}, {"name": "remove_worktree", "description": "Remove a worktree. Refuses if uncommitted changes unless discard_changes=true.", "input_schema": {"type": "object", "properties": {"name": {"type": "string"}, "discard_changes": {"type": "boolean"}}, "required": ["name"]}}, {"name": "keep_worktree", "description": "Keep a worktree for manual review.", "input_schema": {"type": "object", "properties": {"name": {"type": "string"}}, "required": ["name"]}}, ] TOOL_HANDLERS = { "bash": run_bash, "read_file": run_read, "write_file": run_write, "create_task": run_create_task, "list_tasks": run_list_tasks, "get_task": run_get_task, "claim_task": run_claim_task, "complete_task": run_complete_task, "spawn_teammate": run_spawn_teammate, "send_message": run_send_message, "request_shutdown": run_request_shutdown, "request_plan": run_request_plan, "review_plan": run_review_plan, "create_worktree": run_create_worktree, "remove_worktree": run_remove_worktree, "keep_worktree": run_keep_worktree, } # ── Context ── MEMORY_DIR = WORKDIR / ".memory" MEMORY_INDEX = MEMORY_DIR / "MEMORY.md" def update_context(context: dict, messages: list) -> dict: memories = "" if MEMORY_INDEX.exists(): memories = MEMORY_INDEX.read_text()[:2000] return {"memories": memories} # ── Agent Loop ── def agent_loop(messages: list, context: dict): system = get_system_prompt(context) while True: try: response = client.messages.create( model=MODEL, system=system, messages=messages, tools=TOOLS, max_tokens=8000) except Exception as e: messages.append({"role": "assistant", "content": [ {"type": "text", "text": f"[Error] {type(e).__name__}: {e}"}]}) return messages.append({"role": "assistant", "content": response.content}) if response.stop_reason != "tool_use": return results = [] for block in response.content: if block.type != "tool_use": continue print(f"\033[36m> {block.name}\033[0m") handler = TOOL_HANDLERS.get(block.name) output = handler(**block.input) if handler else "Unknown" print(str(output)[:300]) results.append({"type": "tool_result", "tool_use_id": block.id, "content": output}) messages.append({"role": "user", "content": results}) context = update_context(context, messages) system = get_system_prompt(context) if __name__ == "__main__": print("s17: worktree isolation") print("Enter a question, press Enter to send. Type q to quit.\n") history = [] context = {"memories": ""} events = queue.Queue() def input_reader(): while True: try: line = input("\033[36ms17 >> \033[0m") except (EOFError, KeyboardInterrupt): events.put(("quit", None)) return events.put(("user", line)) def inbox_poller(): while True: time.sleep(1) if BUS.peek("lead"): events.put(("wake", None)) threading.Thread(target=input_reader, daemon=True).start() threading.Thread(target=inbox_poller, daemon=True).start() had_teammates = False while True: kind, payload = events.get() if kind == "quit": break if kind == "user": if payload.strip().lower() in ("q", "exit", ""): break history.append({"role": "user", "content": payload}) else: inbox = consume_lead_inbox(route_protocol=True) if not inbox: continue history.append({"role": "user", "content": format_team_events(inbox)}) print(f"\n\033[33m[wake: {len(inbox)} team events " f"-> new turn]\033[0m") agent_loop(history, context) context = update_context(context, history) for block in history[-1]["content"]: if getattr(block, "type", None) == "text": print(block.text) elif isinstance(block, dict) and block.get("type") == "text": print(block.get("text", "")) if active_teammates: had_teammates = True elif had_teammates and not BUS.peek("lead"): print("\033[32m[all teammates shut down]\033[0m") had_teammates = False print()