feat: add MCP server control plane for Auto Claude pipeline

Implements a FastMCP-based MCP server that exposes the full Auto Claude
backend as 43 tools across 12 categories. Any MCP client (Claude Code,
Claude Desktop, future web app) can now manage the autonomous coding
pipeline programmatically.

Architecture:
- FastMCP server with stdio/SSE/HTTP transport options
- Service layer wrapping existing backend modules (no duplication)
- Long-running operation tracker with polling pattern
- Stdout isolation to prevent MCP protocol corruption

Tool categories (43 total):
- Project (4): set_active, get_status, list_specs, get_index
- Tasks (6): list, create, get, update, delete, update_status
- Specs (4): create, get_status, get_content, list
- Execution (4): build_start/stop/progress/logs
- QA (3): start_review, get_report, approve
- Workspace (5): list, diff, merge, discard, create_pr
- GitHub (5): review_pr, list_issues, auto_fix, get_review, triage
- Insights (2): ask, suggest_tasks
- Roadmap (3): generate, get, refresh
- Ideation (2): generate, get
- Memory (3): search, add_episode, get_recent
- Operations (2): get_status, cancel

Usage:
  python -m mcp_server --project-dir /path/to/project

Co-Authored-By: Claude Opus 4.6 <[email protected]>
This commit is contained in:
AndyMik90
2026-02-15 20:23:11 +01:00
co-authored by Claude Opus 4.6
parent 2e4b5ac659
commit 8ab939e11e
29 changed files with 4072 additions and 0 deletions
+9
View File
@@ -0,0 +1,9 @@
"""
Auto Claude MCP Server
======================
Control plane for the Auto Claude autonomous coding pipeline.
Exposes all backend capabilities via the Model Context Protocol (MCP).
"""
__version__ = "0.1.0"
+95
View File
@@ -0,0 +1,95 @@
"""
Auto Claude MCP Server Entry Point
===================================
Usage:
python -m mcp_server --project-dir /path/to/project
python -m mcp_server --project-dir /path/to/project --transport sse --port 8642
"""
from __future__ import annotations
import argparse
import logging
import sys
# Configure logging to stderr (stdout is reserved for MCP protocol over stdio)
logging.basicConfig(
level=logging.INFO,
format="%(asctime)s [%(name)s] %(levelname)s: %(message)s",
stream=sys.stderr,
)
logger = logging.getLogger("mcp_server")
def main() -> None:
parser = argparse.ArgumentParser(
description="Auto Claude MCP Server - control plane for the autonomous coding pipeline",
)
parser.add_argument(
"--project-dir",
required=True,
help="Path to the project directory to manage",
)
parser.add_argument(
"--transport",
choices=["stdio", "sse", "streamable-http"],
default="stdio",
help="MCP transport to use (default: stdio)",
)
parser.add_argument(
"--port",
type=int,
default=8642,
help="Port for SSE/HTTP transport (default: 8642)",
)
parser.add_argument(
"--host",
default="127.0.0.1",
help="Host for SSE/HTTP transport (default: 127.0.0.1)",
)
parser.add_argument(
"--debug",
action="store_true",
help="Enable debug logging",
)
args = parser.parse_args()
if args.debug:
logging.getLogger().setLevel(logging.DEBUG)
# Initialize project context (adds backend to sys.path, loads .env)
from mcp_server.config import initialize
initialize(args.project_dir)
# Import server and register tools AFTER initialization
# (tools need backend modules on sys.path)
from mcp_server.server import mcp, register_all_tools
register_all_tools()
logger.info(
"Starting Auto Claude MCP server (transport=%s, project=%s)",
args.transport,
args.project_dir,
)
# For stdio transport, redirect any stray stdout prints to stderr
# to prevent corrupting the MCP JSON-RPC protocol
if args.transport == "stdio":
# Capture any prints from backend modules that write to stdout
_original_stdout = sys.stdout
sys.stdout = sys.stderr
# Run the server
if args.transport == "stdio":
mcp.run(transport="stdio")
elif args.transport == "sse":
mcp.run(transport="sse", host=args.host, port=args.port)
elif args.transport == "streamable-http":
mcp.run(transport="streamable-http", host=args.host, port=args.port)
if __name__ == "__main__":
main()
+114
View File
@@ -0,0 +1,114 @@
"""
MCP Server Configuration
========================
Manages project context and backend initialization for the MCP server.
The project directory is set once at startup and used by all tools.
"""
from __future__ import annotations
import json
import logging
import sys
from pathlib import Path
logger = logging.getLogger(__name__)
# Global project context - set once at server startup
_project_dir: Path | None = None
_auto_claude_dir: Path | None = None
def initialize(project_dir: str | Path) -> None:
"""Initialize the MCP server with a project directory.
This sets up the Python path so backend modules can be imported,
loads the .env file, and validates the project structure.
Args:
project_dir: Path to the user's project directory
"""
global _project_dir, _auto_claude_dir
_project_dir = Path(project_dir).resolve()
if not _project_dir.is_dir():
raise ValueError(f"Project directory does not exist: {_project_dir}")
# Add backend to sys.path so existing modules can be imported
backend_dir = Path(__file__).parent.parent
if str(backend_dir) not in sys.path:
sys.path.insert(0, str(backend_dir))
# Load .env if present
try:
from cli.utils import import_dotenv
load_dotenv = import_dotenv()
env_file = backend_dir / ".env"
if env_file.exists():
load_dotenv(env_file)
except Exception:
logger.debug("Could not load .env file (non-critical)")
# Determine .auto-claude directory
auto_claude = _project_dir / ".auto-claude"
if not auto_claude.is_dir():
# Also check legacy 'auto-claude' (no dot prefix)
alt = _project_dir / "auto-claude"
if alt.is_dir():
auto_claude = alt
else:
logger.warning(
"No .auto-claude directory found in %s. "
"Some tools may not work until the project is initialized.",
_project_dir,
)
_auto_claude_dir = auto_claude
logger.info("MCP server initialized for project: %s", _project_dir)
def get_project_dir() -> Path:
"""Get the active project directory. Raises if not initialized."""
if _project_dir is None:
raise RuntimeError(
"MCP server not initialized. Call config.initialize() first."
)
return _project_dir
def get_auto_claude_dir() -> Path:
"""Get the .auto-claude directory for the active project."""
if _auto_claude_dir is None:
raise RuntimeError(
"MCP server not initialized. Call config.initialize() first."
)
return _auto_claude_dir
def get_specs_dir() -> Path:
"""Get the specs directory for the active project."""
return get_auto_claude_dir() / "specs"
def get_project_index() -> dict:
"""Load and return the project index if available."""
index_path = get_auto_claude_dir() / "project_index.json"
if not index_path.exists():
return {}
try:
with open(index_path, encoding="utf-8") as f:
return json.load(f)
except (json.JSONDecodeError, OSError) as e:
logger.warning("Failed to load project index: %s", e)
return {}
def is_initialized() -> bool:
"""Check if the project has been initialized with .auto-claude."""
try:
ac_dir = get_auto_claude_dir()
return ac_dir.is_dir()
except RuntimeError:
return False
+158
View File
@@ -0,0 +1,158 @@
"""
Long-Running Operation Tracker
===============================
Tracks async operations (spec creation, builds, QA, etc.) so MCP clients
can poll for progress. Tools that start long-running work return an
operation_id immediately; clients poll operation_get_status() for updates.
"""
from __future__ import annotations
import asyncio
import logging
import time
import uuid
from dataclasses import dataclass, field
from enum import Enum
from typing import Any
logger = logging.getLogger(__name__)
class OperationStatus(str, Enum):
PENDING = "pending"
RUNNING = "running"
COMPLETED = "completed"
FAILED = "failed"
CANCELLED = "cancelled"
@dataclass
class Operation:
"""Represents a long-running operation."""
id: str
type: str # e.g. "spec_create", "build", "qa_review"
status: OperationStatus = OperationStatus.PENDING
progress: int = 0 # 0-100
message: str = ""
result: Any = None
error: str | None = None
created_at: float = field(default_factory=time.time)
updated_at: float = field(default_factory=time.time)
_task: asyncio.Task | None = field(default=None, repr=False)
def to_dict(self) -> dict:
"""Serialize for MCP response."""
return {
"id": self.id,
"type": self.type,
"status": self.status.value,
"progress": self.progress,
"message": self.message,
"result": self.result,
"error": self.error,
"created_at": self.created_at,
"updated_at": self.updated_at,
"elapsed_seconds": round(time.time() - self.created_at, 1),
}
class OperationTracker:
"""Manages the lifecycle of long-running operations."""
def __init__(self, max_completed: int = 100):
self._operations: dict[str, Operation] = {}
self._max_completed = max_completed
def create(self, operation_type: str, message: str = "") -> Operation:
"""Create a new operation and return it."""
op = Operation(
id=str(uuid.uuid4()),
type=operation_type,
status=OperationStatus.PENDING,
message=message or f"Starting {operation_type}...",
)
self._operations[op.id] = op
self._cleanup_old()
return op
def get(self, operation_id: str) -> Operation | None:
"""Get an operation by ID."""
return self._operations.get(operation_id)
def update(
self,
operation_id: str,
*,
status: OperationStatus | None = None,
progress: int | None = None,
message: str | None = None,
result: Any = None,
error: str | None = None,
) -> Operation | None:
"""Update an operation's state."""
op = self._operations.get(operation_id)
if op is None:
return None
if status is not None:
op.status = status
if progress is not None:
op.progress = max(0, min(100, progress))
if message is not None:
op.message = message
if result is not None:
op.result = result
if error is not None:
op.error = error
op.updated_at = time.time()
return op
def cancel(self, operation_id: str) -> bool:
"""Cancel a running operation."""
op = self._operations.get(operation_id)
if op is None:
return False
if op.status in (OperationStatus.COMPLETED, OperationStatus.FAILED):
return False
# Cancel the asyncio task if it exists
if op._task and not op._task.done():
op._task.cancel()
op.status = OperationStatus.CANCELLED
op.message = "Operation cancelled by user"
op.updated_at = time.time()
return True
def list_active(self) -> list[Operation]:
"""List all active (non-terminal) operations."""
return [
op
for op in self._operations.values()
if op.status in (OperationStatus.PENDING, OperationStatus.RUNNING)
]
def _cleanup_old(self) -> None:
"""Remove old completed operations to prevent memory growth."""
completed = [
op
for op in self._operations.values()
if op.status
in (
OperationStatus.COMPLETED,
OperationStatus.FAILED,
OperationStatus.CANCELLED,
)
]
if len(completed) > self._max_completed:
# Sort by created_at, remove oldest
completed.sort(key=lambda o: o.created_at)
for op in completed[: len(completed) - self._max_completed]:
del self._operations[op.id]
# Global singleton
tracker = OperationTracker()
+56
View File
@@ -0,0 +1,56 @@
"""
Auto Claude MCP Server
======================
FastMCP server instance with all tool registrations.
Tools are organized into modules under mcp_server/tools/.
Each module's register() function adds tools to the server.
"""
from __future__ import annotations
import logging
from fastmcp import FastMCP
logger = logging.getLogger(__name__)
# Create the FastMCP server instance
mcp = FastMCP(
"Auto Claude",
instructions=(
"Auto Claude is an autonomous multi-agent coding framework. "
"Use these tools to manage tasks, create specs, run builds, "
"perform QA reviews, manage workspaces, and more. "
"Long-running operations return an operation_id - "
"poll with operation_get_status() for progress."
),
)
def register_all_tools() -> None:
"""Register all tool modules with the MCP server.
Each tool module defines functions decorated with @mcp.tool()
that are imported here to trigger registration.
"""
# Phase 1: Project & Task management
# Phase 2: Core autonomous pipeline
# Phase 3: Feature tools
# Operations management (poll long-running ops)
from mcp_server.tools import (
execution, # noqa: F401
github, # noqa: F401
ideation, # noqa: F401
insights, # noqa: F401
memory, # noqa: F401
ops, # noqa: F401
project, # noqa: F401
qa, # noqa: F401
roadmap, # noqa: F401
specs, # noqa: F401
tasks, # noqa: F401
workspace, # noqa: F401
)
logger.info("All MCP tools registered successfully")
@@ -0,0 +1 @@
"""Service layer - thin adapters wrapping existing backend modules."""
@@ -0,0 +1,322 @@
"""
Execution Service
==================
Service layer for spawning and managing build processes.
Wraps the run.py subprocess and parses task events from stdout.
"""
from __future__ import annotations
import asyncio
import json
import logging
import sys
from pathlib import Path
logger = logging.getLogger(__name__)
# Matches core/task_event.py
TASK_EVENT_PREFIX = "__TASK_EVENT__:"
class ExecutionService:
"""Manages build execution as a subprocess of run.py."""
def __init__(self, project_dir: Path):
self.project_dir = project_dir
self._processes: dict[str, asyncio.subprocess.Process] = {}
self._logs: dict[str, list[str]] = {}
self._events: dict[str, list[dict]] = {}
async def start_build(
self,
spec_id: str,
model: str = "sonnet",
thinking_level: str = "medium",
) -> asyncio.subprocess.Process:
"""Spawn a build subprocess for the given spec.
Args:
spec_id: The spec folder name
model: Model shorthand
thinking_level: Thinking level
Returns:
The subprocess handle
Raises:
RuntimeError: If a build is already running for this spec
"""
if spec_id in self._processes:
proc = self._processes[spec_id]
if proc.returncode is None:
raise RuntimeError(
f"Build already running for spec '{spec_id}'. "
"Stop it first with build_stop()."
)
backend_dir = Path(__file__).parent.parent.parent # apps/backend/
run_py = backend_dir / "run.py"
if not run_py.exists():
raise FileNotFoundError(f"run.py not found at {run_py}")
cmd = [
sys.executable,
str(run_py),
"--spec",
spec_id,
"--project-dir",
str(self.project_dir),
"--model",
model,
"--thinking",
thinking_level,
]
proc = await asyncio.create_subprocess_exec(
*cmd,
stdout=asyncio.subprocess.PIPE,
stderr=asyncio.subprocess.PIPE,
cwd=str(backend_dir),
)
self._processes[spec_id] = proc
self._logs[spec_id] = []
self._events[spec_id] = []
# Start background reader for stdout
asyncio.create_task(self._read_output(spec_id, proc))
return proc
async def _read_output(
self, spec_id: str, proc: asyncio.subprocess.Process
) -> None:
"""Read stdout from the build process, parsing task events.
Args:
spec_id: The spec being built
proc: The subprocess to read from
"""
if proc.stdout is None:
return
try:
while True:
line_bytes = await proc.stdout.readline()
if not line_bytes:
break
line = line_bytes.decode("utf-8", errors="replace").rstrip("\n")
# Store the log line
log_list = self._logs.get(spec_id)
if log_list is not None:
log_list.append(line)
# Cap stored logs to prevent unbounded growth
if len(log_list) > 5000:
del log_list[:1000]
# Parse task events
event = self.parse_event(line)
if event is not None:
events_list = self._events.get(spec_id)
if events_list is not None:
events_list.append(event)
except Exception as e:
logger.warning("Error reading build output for %s: %s", spec_id, e)
def parse_event(self, line: str) -> dict | None:
"""Parse a task event line from build stdout.
Args:
line: A line of stdout output
Returns:
Parsed event dict or None if not an event line
"""
if not line.startswith(TASK_EVENT_PREFIX):
return None
try:
return json.loads(line[len(TASK_EVENT_PREFIX) :])
except (json.JSONDecodeError, ValueError):
return None
def stop_build(self, spec_id: str) -> dict:
"""Stop a running build process.
Args:
spec_id: The spec being built
Returns:
Status dict
"""
proc = self._processes.get(spec_id)
if proc is None:
return {"success": False, "error": f"No build found for spec '{spec_id}'"}
if proc.returncode is not None:
return {
"success": False,
"error": f"Build for '{spec_id}' already finished (exit code {proc.returncode})",
}
try:
proc.terminate()
return {"success": True, "message": f"Build for '{spec_id}' terminated"}
except ProcessLookupError:
return {"success": False, "error": "Process already exited"}
def get_progress(self, spec_id: str) -> dict:
"""Get progress of a build by inspecting events and process state.
Args:
spec_id: The spec being built
Returns:
Dict with status, events, and process info
"""
proc = self._processes.get(spec_id)
events = self._events.get(spec_id, [])
if proc is None:
# Check if there's a completed implementation plan on disk
return self._get_disk_progress(spec_id)
is_running = proc.returncode is None
latest_event = events[-1] if events else None
return {
"spec_id": spec_id,
"running": is_running,
"exit_code": proc.returncode,
"event_count": len(events),
"latest_event": latest_event,
"log_lines": len(self._logs.get(spec_id, [])),
}
def get_logs(self, spec_id: str, tail: int = 50) -> dict:
"""Get recent build logs.
Args:
spec_id: The spec being built
tail: Number of recent lines to return
Returns:
Dict with log lines
"""
logs = self._logs.get(spec_id, [])
if not logs:
# Try to find logs on disk
return self._get_disk_logs(spec_id, tail)
return {
"spec_id": spec_id,
"total_lines": len(logs),
"lines": logs[-tail:],
}
def _get_disk_progress(self, spec_id: str) -> dict:
"""Check on-disk state for build progress when no process is tracked.
Args:
spec_id: The spec folder name
Returns:
Progress dict from disk state
"""
specs_dir = self.project_dir / ".auto-claude" / "specs"
spec_dir = self._resolve_spec_dir(specs_dir, spec_id)
if spec_dir is None:
return {"error": f"Spec '{spec_id}' not found"}
plan_file = spec_dir / "implementation_plan.json"
if not plan_file.exists():
return {
"spec_id": spec_id,
"running": False,
"status": "no_plan",
"message": "No implementation plan found. Create a spec first.",
}
try:
with open(plan_file, encoding="utf-8") as f:
plan = json.load(f)
except (json.JSONDecodeError, OSError):
return {
"spec_id": spec_id,
"running": False,
"status": "error",
"message": "Could not read implementation plan",
}
subtasks = plan.get("subtasks", [])
completed = sum(1 for s in subtasks if s.get("status") == "completed")
total = len(subtasks)
qa_signoff = plan.get("qa_signoff")
if qa_signoff and qa_signoff.get("status") == "approved":
status = "qa_approved"
elif qa_signoff and qa_signoff.get("status") == "rejected":
status = "qa_rejected"
elif completed == total and total > 0:
status = "build_complete"
elif completed > 0:
status = "building"
else:
status = "not_started"
return {
"spec_id": spec_id,
"running": False,
"status": status,
"subtasks_completed": completed,
"subtasks_total": total,
"qa_signoff": qa_signoff,
}
def _get_disk_logs(self, spec_id: str, tail: int) -> dict:
"""Try to find build logs on disk.
Args:
spec_id: The spec folder name
tail: Number of lines to return
Returns:
Dict with log content
"""
specs_dir = self.project_dir / ".auto-claude" / "specs"
spec_dir = self._resolve_spec_dir(specs_dir, spec_id)
if spec_dir is None:
return {"error": f"Spec '{spec_id}' not found", "lines": []}
# Check for task log file
log_file = spec_dir / "task_log.jsonl"
if not log_file.exists():
return {
"spec_id": spec_id,
"lines": [],
"message": "No build logs found",
}
try:
lines = log_file.read_text(encoding="utf-8").strip().split("\n")
return {
"spec_id": spec_id,
"total_lines": len(lines),
"lines": lines[-tail:],
}
except OSError as e:
return {"error": str(e), "lines": []}
def _resolve_spec_dir(self, specs_dir: Path, spec_id: str) -> Path | None:
"""Resolve spec_id to directory with prefix matching."""
exact = specs_dir / spec_id
if exact.is_dir():
return exact
if specs_dir.is_dir():
for item in specs_dir.iterdir():
if item.is_dir() and item.name.startswith(spec_id):
return item
return None
@@ -0,0 +1,202 @@
"""
GitHub Service
==============
Wraps the backend GitHubOrchestrator for MCP tool access.
Handles repo detection, config creation, and result serialization.
"""
from __future__ import annotations
import json
import logging
import subprocess
from pathlib import Path
logger = logging.getLogger(__name__)
class GitHubService:
"""Service layer for GitHub automation features."""
def __init__(self, project_dir: Path):
self.project_dir = project_dir
self.github_dir = project_dir / ".auto-claude" / "github"
def _detect_repo(self) -> str | None:
"""Detect owner/repo from git remote origin."""
try:
result = subprocess.run(
["git", "remote", "get-url", "origin"],
capture_output=True,
text=True,
cwd=str(self.project_dir),
timeout=10,
)
if result.returncode != 0:
return None
url = result.stdout.strip()
# Handle SSH: [email protected]:owner/repo.git
if url.startswith("git@"):
parts = url.split(":")[-1]
return parts.removesuffix(".git")
# Handle HTTPS: https://github.com/owner/repo.git
if "github.com" in url:
parts = url.split("github.com/")[-1]
return parts.removesuffix(".git")
return None
except Exception as e:
logger.warning("Failed to detect repo from git remote: %s", e)
return None
def _get_repo(self, repo: str | None) -> str:
"""Get repo string, falling back to auto-detection."""
if repo:
return repo
detected = self._detect_repo()
if not detected:
raise ValueError(
"Could not detect repository. Provide 'repo' parameter "
"in owner/repo format, or ensure a GitHub remote is configured."
)
return detected
def _create_config(self, repo: str, model: str = "sonnet"):
"""Create a GitHubRunnerConfig with sensible defaults."""
# Get GitHub token from environment
import os
from runners.github.models import GitHubRunnerConfig
token = os.environ.get("GITHUB_TOKEN", "")
if not token:
# Try gh CLI auth token
try:
result = subprocess.run(
["gh", "auth", "token"],
capture_output=True,
text=True,
timeout=10,
)
if result.returncode == 0:
token = result.stdout.strip()
except Exception:
pass
return GitHubRunnerConfig(
token=token,
repo=repo,
model=model,
thinking_level="medium",
pr_review_enabled=True,
triage_enabled=True,
)
async def review_pr(
self, pr_number: int, repo: str | None = None, model: str = "sonnet"
) -> dict:
"""Review a pull request with AI."""
try:
from runners.github.orchestrator import GitHubOrchestrator
resolved_repo = self._get_repo(repo)
config = self._create_config(resolved_repo, model)
orchestrator = GitHubOrchestrator(
project_dir=self.project_dir, config=config
)
result = await orchestrator.review_pr(pr_number)
return {"success": True, "data": result.to_dict()}
except ImportError:
return {"error": "GitHub runner module not available"}
except Exception as e:
return {"error": str(e)}
async def list_issues(
self, state: str = "open", limit: int = 30, repo: str | None = None
) -> dict:
"""List GitHub issues using gh CLI."""
try:
resolved_repo = self._get_repo(repo)
cmd = [
"gh",
"issue",
"list",
"--repo",
resolved_repo,
"--state",
state,
"--limit",
str(limit),
"--json",
"number,title,state,labels,author,createdAt,updatedAt",
]
result = subprocess.run(
cmd,
capture_output=True,
text=True,
cwd=str(self.project_dir),
timeout=30,
)
if result.returncode != 0:
return {"error": f"gh CLI failed: {result.stderr.strip()}"}
issues = json.loads(result.stdout)
return {"success": True, "issues": issues, "count": len(issues)}
except Exception as e:
return {"error": str(e)}
async def auto_fix_issue(self, issue_number: int, repo: str | None = None) -> dict:
"""Auto-fix a GitHub issue."""
try:
from runners.github.orchestrator import GitHubOrchestrator
resolved_repo = self._get_repo(repo)
config = self._create_config(resolved_repo)
config.auto_fix_enabled = True
orchestrator = GitHubOrchestrator(
project_dir=self.project_dir, config=config
)
state = await orchestrator.auto_fix_issue(issue_number)
return {"success": True, "data": state.to_dict()}
except ImportError:
return {"error": "GitHub runner module not available"}
except Exception as e:
return {"error": str(e)}
def get_review(self, pr_number: int) -> dict:
"""Get the most recent review result for a PR."""
try:
from runners.github.models import PRReviewResult
result = PRReviewResult.load(self.github_dir, pr_number)
if result is None:
return {"error": f"No review found for PR #{pr_number}"}
return {"success": True, "data": result.to_dict()}
except ImportError:
return {"error": "GitHub runner module not available"}
except Exception as e:
return {"error": str(e)}
async def triage_issues(
self, issue_numbers: list[int], repo: str | None = None
) -> dict:
"""Triage and classify GitHub issues."""
try:
from runners.github.orchestrator import GitHubOrchestrator
resolved_repo = self._get_repo(repo)
config = self._create_config(resolved_repo)
config.triage_enabled = True
orchestrator = GitHubOrchestrator(
project_dir=self.project_dir, config=config
)
results = await orchestrator.triage_issues(issue_numbers=issue_numbers)
return {
"success": True,
"data": [r.to_dict() for r in results],
"count": len(results),
}
except ImportError:
return {"error": "GitHub runner module not available"}
except Exception as e:
return {"error": str(e)}
@@ -0,0 +1,82 @@
"""
Ideation Service
=================
Wraps the backend IdeationOrchestrator for MCP tool access.
"""
from __future__ import annotations
import json
import logging
from pathlib import Path
logger = logging.getLogger(__name__)
# Valid ideation types
VALID_IDEATION_TYPES = [
"low_hanging_fruit",
"ui_ux_improvements",
"high_value_features",
]
class IdeationService:
"""Service layer for AI-powered ideation generation."""
def __init__(self, project_dir: Path):
self.project_dir = project_dir
self.ideation_dir = project_dir / ".auto-claude" / "ideation"
async def generate(
self,
types: list[str] | None = None,
refresh: bool = False,
model: str = "sonnet",
thinking_level: str = "medium",
) -> dict:
"""Generate ideas for project improvements."""
try:
from ideation import IdeationOrchestrator
# Validate types
enabled_types = types or VALID_IDEATION_TYPES
invalid = [t for t in enabled_types if t not in VALID_IDEATION_TYPES]
if invalid:
return {
"error": f"Invalid ideation types: {invalid}. "
f"Valid types: {VALID_IDEATION_TYPES}"
}
orchestrator = IdeationOrchestrator(
project_dir=self.project_dir,
enabled_types=enabled_types,
model=model,
thinking_level=thinking_level,
refresh=refresh,
)
success = await orchestrator.run()
if success:
return self.get_ideation()
return {"error": "Ideation generation failed. Check logs for details."}
except ImportError:
return {"error": "Ideation module not available"}
except Exception as e:
return {"error": str(e)}
def get_ideation(self) -> dict:
"""Get previously generated ideation results from disk."""
ideation_file = self.ideation_dir / "ideation.json"
if not ideation_file.exists():
return {
"success": True,
"data": None,
"message": "No ideation data yet. Use ideation_generate first.",
}
try:
with open(ideation_file, encoding="utf-8") as f:
ideation = json.load(f)
return {"success": True, "data": ideation}
except (json.JSONDecodeError, OSError) as e:
return {"error": f"Failed to load ideation data: {e}"}
@@ -0,0 +1,116 @@
"""
Insights Service
=================
Wraps the backend InsightsRunner for MCP tool access.
Captures stdout output since run_with_sdk prints to stdout.
"""
from __future__ import annotations
import contextlib
import io
import json
import logging
from pathlib import Path
logger = logging.getLogger(__name__)
class InsightsService:
"""Service layer for codebase insights / AI chat."""
def __init__(self, project_dir: Path):
self.project_dir = project_dir
async def ask(
self,
question: str,
history: list | None = None,
model: str = "sonnet",
thinking_level: str = "medium",
) -> dict:
"""Ask an AI question about the codebase.
IMPORTANT: run_with_sdk prints to stdout, so we capture it.
"""
try:
from runners.insights_runner import run_with_sdk
history = history or []
captured = io.StringIO()
with contextlib.redirect_stdout(captured):
await run_with_sdk(
project_dir=str(self.project_dir),
message=question,
history=history,
model=model,
thinking_level=thinking_level,
)
output = captured.getvalue()
# Parse out any task suggestions from the output
task_suggestions = []
response_lines = []
for line in output.split("\n"):
if line.startswith("__TASK_SUGGESTION__:"):
try:
suggestion_json = line.split("__TASK_SUGGESTION__:", 1)[1]
task_suggestions.append(json.loads(suggestion_json))
except (json.JSONDecodeError, IndexError):
pass
elif line.startswith("__TOOL_START__:") or line.startswith(
"__TOOL_END__:"
):
# Skip tool markers - they're for the Electron UI
pass
else:
response_lines.append(line)
response_text = "\n".join(response_lines).strip()
return {
"success": True,
"response": response_text,
"task_suggestions": task_suggestions,
}
except ImportError:
return {"error": "Insights runner module not available"}
except Exception as e:
return {"error": str(e)}
def suggest_tasks(self) -> dict:
"""Get AI-suggested tasks based on project analysis.
Reads the most recent ideation/insights data if available.
"""
try:
ideation_file = (
self.project_dir / ".auto-claude" / "ideation" / "ideation.json"
)
if ideation_file.exists():
with open(ideation_file, encoding="utf-8") as f:
ideation = json.load(f)
ideas = ideation.get("ideas", [])
# Convert top ideas to task suggestions
suggestions = []
for idea in ideas[:10]:
suggestions.append(
{
"title": idea.get("title", ""),
"description": idea.get("description", ""),
"category": idea.get("type", "feature"),
"impact": idea.get("impact", "medium"),
"effort": idea.get("effort", "medium"),
}
)
return {"success": True, "suggestions": suggestions}
return {
"success": True,
"suggestions": [],
"message": "No ideation data available. Run ideation_generate first.",
}
except Exception as e:
return {"error": str(e)}
@@ -0,0 +1,146 @@
"""
Memory Service
===============
Wraps the Graphiti memory system for MCP tool access.
Gracefully handles the case where Graphiti is not enabled/configured.
"""
from __future__ import annotations
import logging
import os
from pathlib import Path
logger = logging.getLogger(__name__)
def _is_graphiti_enabled() -> bool:
"""Check if Graphiti memory is enabled via environment variable."""
return os.environ.get("GRAPHITI_ENABLED", "").lower() in ("true", "1")
class MemoryService:
"""Service layer for Graphiti-based semantic memory."""
def __init__(self, project_dir: Path):
self.project_dir = project_dir
self._memory = None
def _get_disabled_message(self) -> dict:
"""Return a helpful error when Graphiti is not enabled."""
return {
"error": "Graphiti memory is not enabled. "
"Set GRAPHITI_ENABLED=true in your .env file and configure "
"the required provider settings (LLM and embedder). "
"See the project documentation for setup instructions."
}
async def _get_memory(self):
"""Lazily initialize and return a GraphitiMemory instance."""
if self._memory is not None:
return self._memory
if not _is_graphiti_enabled():
return None
try:
from integrations.graphiti.memory import (
GraphitiMemory,
GroupIdMode,
)
# Use a dummy spec_dir since we're in project-wide mode
spec_dir = self.project_dir / ".auto-claude" / "mcp_memory"
spec_dir.mkdir(parents=True, exist_ok=True)
memory = GraphitiMemory(
spec_dir=spec_dir,
project_dir=self.project_dir,
group_id_mode=GroupIdMode.PROJECT,
)
if not await memory.initialize():
logger.warning("Failed to initialize Graphiti memory")
return None
self._memory = memory
return memory
except ImportError:
logger.warning("Graphiti modules not available")
return None
except Exception as e:
logger.warning("Failed to create Graphiti memory: %s", e)
return None
async def search(self, query: str, limit: int = 10) -> dict:
"""Search the project's semantic memory."""
if not _is_graphiti_enabled():
return self._get_disabled_message()
try:
memory = await self._get_memory()
if memory is None:
return {"error": "Could not initialize Graphiti memory"}
results = await memory._search.get_relevant_context(
query=query,
num_results=limit,
)
return {
"success": True,
"results": results,
"count": len(results),
}
except Exception as e:
return {"error": f"Memory search failed: {e}"}
async def add_episode(self, content: str, source: str = "mcp") -> dict:
"""Add a new episode/fact to the project's memory."""
if not _is_graphiti_enabled():
return self._get_disabled_message()
try:
memory = await self._get_memory()
if memory is None:
return {"error": "Could not initialize Graphiti memory"}
success = await memory.save_session_insights(
session_num=0,
insights={
"content": content,
"source": source,
"type": "mcp_episode",
},
)
if success:
return {"success": True, "message": "Episode added to memory"}
return {"error": "Failed to save episode to memory"}
except Exception as e:
return {"error": f"Failed to add episode: {e}"}
async def get_recent(self, limit: int = 10) -> dict:
"""Get recent memory entries."""
if not _is_graphiti_enabled():
return self._get_disabled_message()
try:
memory = await self._get_memory()
if memory is None:
return {"error": "Could not initialize Graphiti memory"}
# Use a broad search to get recent entries
results = await memory._search.get_relevant_context(
query="recent project activity and insights",
num_results=limit,
)
return {
"success": True,
"results": results,
"count": len(results),
}
except Exception as e:
return {"error": f"Failed to get recent memory: {e}"}
@@ -0,0 +1,236 @@
"""
QA Service
===========
Service layer wrapping the backend QA reviewer for MCP tool consumption.
Handles client creation, stdout isolation, and error management.
"""
from __future__ import annotations
import contextlib
import io
import json
import logging
from pathlib import Path
logger = logging.getLogger(__name__)
class QAService:
"""Wraps QA review and approval operations for MCP server use."""
def __init__(self, project_dir: Path):
self.project_dir = project_dir
async def start_review(
self,
spec_id: str,
model: str = "sonnet",
thinking_level: str = "medium",
max_iterations: int = 3,
) -> dict:
"""Run a QA review session for a completed build.
Args:
spec_id: The spec folder name
model: Model shorthand
thinking_level: Thinking level
max_iterations: Maximum QA loop iterations
Returns:
Dict with review outcome (approved/rejected/error)
"""
spec_dir = self._resolve_spec_dir(spec_id)
if spec_dir is None:
return {"error": f"Spec '{spec_id}' not found"}
# Verify the build is complete before starting QA
plan_file = spec_dir / "implementation_plan.json"
if not plan_file.exists():
return {
"error": "No implementation plan found. Build the spec first.",
}
try:
from core.client import create_client
from qa.reviewer import run_qa_agent_session
except ImportError as e:
logger.error("Failed to import QA modules: %s", e)
return {"error": f"Backend module not available: {e}"}
try:
# Determine QA session number from existing state
qa_session = self._get_next_qa_session(spec_dir)
# Create a Claude SDK client for the QA agent
captured = io.StringIO()
with contextlib.redirect_stdout(captured):
client = create_client(
project_dir=self.project_dir,
spec_dir=spec_dir,
model=model,
phase="qa_reviewer",
)
status, response_text, error_info = await run_qa_agent_session(
client=client,
project_dir=self.project_dir,
spec_dir=spec_dir,
qa_session=qa_session,
max_iterations=max_iterations,
)
return {
"spec_id": spec_id,
"status": status,
"qa_session": qa_session,
"response_preview": response_text[:1000] if response_text else "",
"error_info": error_info if error_info else None,
"output": captured.getvalue()[-1000:] if captured.getvalue() else "",
}
except Exception as e:
logger.exception("QA review failed for %s", spec_id)
return {
"spec_id": spec_id,
"status": "error",
"error": str(e),
}
def get_report(self, spec_id: str) -> dict:
"""Get the QA report for a spec.
Args:
spec_id: The spec folder name
Returns:
Dict with QA report content and status
"""
spec_dir = self._resolve_spec_dir(spec_id)
if spec_dir is None:
return {"error": f"Spec '{spec_id}' not found"}
result: dict = {"spec_id": spec_id}
# Read qa_report.md
qa_report = spec_dir / "qa_report.md"
if qa_report.exists():
try:
result["report"] = qa_report.read_text(encoding="utf-8")
except OSError as e:
result["report_error"] = str(e)
# Read QA fix request if present
fix_request = spec_dir / "QA_FIX_REQUEST.md"
if fix_request.exists():
try:
result["fix_request"] = fix_request.read_text(encoding="utf-8")
except OSError as e:
result["fix_request_error"] = str(e)
# Read qa_signoff from implementation plan
plan_file = spec_dir / "implementation_plan.json"
if plan_file.exists():
try:
with open(plan_file, encoding="utf-8") as f:
plan = json.load(f)
qa_signoff = plan.get("qa_signoff")
if qa_signoff:
result["qa_signoff"] = qa_signoff
except (json.JSONDecodeError, OSError):
pass
if "report" not in result and "qa_signoff" not in result:
result["message"] = "No QA report found. Run QA review first."
return result
def approve(self, spec_id: str) -> dict:
"""Manually approve a spec's QA status.
Args:
spec_id: The spec folder name
Returns:
Dict with approval result
"""
spec_dir = self._resolve_spec_dir(spec_id)
if spec_dir is None:
return {"error": f"Spec '{spec_id}' not found"}
plan_file = spec_dir / "implementation_plan.json"
if not plan_file.exists():
return {"error": "No implementation plan found"}
try:
with open(plan_file, encoding="utf-8") as f:
plan = json.load(f)
except (json.JSONDecodeError, OSError) as e:
return {"error": f"Could not read implementation plan: {e}"}
from datetime import datetime, timezone
plan["qa_signoff"] = {
"status": "approved",
"timestamp": datetime.now(timezone.utc).isoformat(),
"qa_session": plan.get("qa_signoff", {}).get("qa_session", 0),
"verified_by": "manual_approval",
"note": "Manually approved via MCP tool",
}
try:
with open(plan_file, "w", encoding="utf-8") as f:
json.dump(plan, f, indent=2)
except OSError as e:
return {"error": f"Could not write implementation plan: {e}"}
return {
"success": True,
"spec_id": spec_id,
"message": "Spec manually approved",
}
def _get_next_qa_session(self, spec_dir: Path) -> int:
"""Get the next QA session number.
Args:
spec_dir: Path to the spec directory
Returns:
Next session number (1-based)
"""
plan_file = spec_dir / "implementation_plan.json"
if not plan_file.exists():
return 1
try:
with open(plan_file, encoding="utf-8") as f:
plan = json.load(f)
qa_signoff = plan.get("qa_signoff", {})
current = qa_signoff.get("qa_session", 0)
return current + 1
except (json.JSONDecodeError, OSError):
return 1
def _resolve_spec_dir(self, spec_id: str) -> Path | None:
"""Resolve spec_id to its directory path.
Args:
spec_id: Full or prefix spec identifier
Returns:
Path to spec directory or None
"""
specs_dir = self.project_dir / ".auto-claude" / "specs"
# Direct match
exact = specs_dir / spec_id
if exact.is_dir():
return exact
# Prefix match
if specs_dir.is_dir():
for item in specs_dir.iterdir():
if item.is_dir() and item.name.startswith(spec_id):
return item
return None
@@ -0,0 +1,65 @@
"""
Roadmap Service
================
Wraps the backend RoadmapOrchestrator for MCP tool access.
"""
from __future__ import annotations
import json
import logging
from pathlib import Path
logger = logging.getLogger(__name__)
class RoadmapService:
"""Service layer for roadmap generation features."""
def __init__(self, project_dir: Path):
self.project_dir = project_dir
self.roadmap_dir = project_dir / ".auto-claude" / "roadmap"
async def generate(
self,
refresh: bool = False,
model: str = "sonnet",
thinking_level: str = "medium",
) -> dict:
"""Generate a strategic roadmap for the project."""
try:
from runners.roadmap.orchestrator import RoadmapOrchestrator
orchestrator = RoadmapOrchestrator(
project_dir=self.project_dir,
model=model,
thinking_level=thinking_level,
refresh=refresh,
)
success = await orchestrator.run()
if success:
# Load and return the generated roadmap
return self.get_roadmap()
return {"error": "Roadmap generation failed. Check logs for details."}
except ImportError:
return {"error": "Roadmap runner module not available"}
except Exception as e:
return {"error": str(e)}
def get_roadmap(self) -> dict:
"""Get the current roadmap data from disk."""
roadmap_file = self.roadmap_dir / "roadmap.json"
if not roadmap_file.exists():
return {
"success": True,
"data": None,
"message": "No roadmap generated yet. Use roadmap_generate first.",
}
try:
with open(roadmap_file, encoding="utf-8") as f:
roadmap = json.load(f)
return {"success": True, "data": roadmap}
except (json.JSONDecodeError, OSError) as e:
return {"error": f"Failed to load roadmap: {e}"}
@@ -0,0 +1,242 @@
"""
Spec Service
=============
Service layer wrapping the backend SpecOrchestrator for MCP tool consumption.
Handles stdout isolation and error management.
"""
from __future__ import annotations
import contextlib
import io
import json
import logging
from pathlib import Path
logger = logging.getLogger(__name__)
class SpecService:
"""Wraps backend spec creation pipeline for MCP server use."""
def __init__(self, project_dir: Path):
self.project_dir = project_dir
async def create_spec(
self,
task_description: str,
model: str = "sonnet",
thinking_level: str = "medium",
complexity_override: str | None = None,
) -> dict:
"""Create a spec using the SpecOrchestrator.
Redirects stdout to prevent protocol corruption when running
under stdio transport.
Args:
task_description: Description of the task to spec out
model: Model shorthand (sonnet, opus, etc.)
thinking_level: Thinking level (low, medium, high)
complexity_override: Force a specific complexity level
Returns:
Dict with success status, spec_dir, spec_id, and any captured output
"""
try:
from spec.pipeline.orchestrator import SpecOrchestrator
except ImportError as e:
logger.error("Failed to import SpecOrchestrator: %s", e)
return {
"success": False,
"error": f"Backend module not available: {e}",
}
try:
captured = io.StringIO()
with contextlib.redirect_stdout(captured):
orchestrator = SpecOrchestrator(
project_dir=self.project_dir,
task_description=task_description,
model=model,
thinking_level=thinking_level,
complexity_override=complexity_override,
use_ai_assessment=True,
)
# Run non-interactively with auto-approve for MCP
success = await orchestrator.run(interactive=False, auto_approve=True)
spec_dir = orchestrator.spec_dir
return {
"success": success,
"spec_dir": str(spec_dir),
"spec_id": spec_dir.name,
"output": captured.getvalue()[-2000:] if captured.getvalue() else "",
}
except Exception as e:
logger.exception("Spec creation failed")
return {
"success": False,
"error": str(e),
}
def get_spec_status(self, spec_id: str) -> dict:
"""Get the status of a spec by checking which phase files exist.
Args:
spec_id: The spec folder name (e.g. '001-my-feature')
Returns:
Dict describing which phases are complete and current state
"""
specs_dir = self.project_dir / ".auto-claude" / "specs"
spec_dir = self._resolve_spec_dir(specs_dir, spec_id)
if spec_dir is None:
return {"error": f"Spec '{spec_id}' not found"}
phases = {
"discovery": (spec_dir / "discovery.md").exists(),
"requirements": (spec_dir / "requirements.json").exists(),
"complexity_assessment": (spec_dir / "complexity_assessment.json").exists(),
"spec": (spec_dir / "spec.md").exists(),
"implementation_plan": (spec_dir / "implementation_plan.json").exists(),
}
# Determine overall status
if phases["implementation_plan"]:
plan = self._load_json(spec_dir / "implementation_plan.json")
qa_signoff = plan.get("qa_signoff") if plan else None
if qa_signoff and qa_signoff.get("status") == "approved":
status = "qa_approved"
elif qa_signoff and qa_signoff.get("status") == "rejected":
status = "qa_rejected"
elif (spec_dir / "qa_report.md").exists():
status = "qa_reviewed"
else:
status = "ready_to_build"
elif phases["spec"]:
status = "spec_complete"
elif phases["requirements"]:
status = "requirements_gathered"
elif phases["discovery"]:
status = "discovery_complete"
else:
status = "pending"
return {
"spec_id": spec_dir.name,
"spec_dir": str(spec_dir),
"status": status,
"phases": phases,
}
def get_spec_content(self, spec_id: str) -> dict:
"""Get the full content of a spec.
Args:
spec_id: The spec folder name
Returns:
Dict with spec.md content, requirements, implementation plan, etc.
"""
specs_dir = self.project_dir / ".auto-claude" / "specs"
spec_dir = self._resolve_spec_dir(specs_dir, spec_id)
if spec_dir is None:
return {"error": f"Spec '{spec_id}' not found"}
content: dict = {
"spec_id": spec_dir.name,
"spec_dir": str(spec_dir),
}
# Read spec.md
spec_md = spec_dir / "spec.md"
if spec_md.exists():
try:
content["spec_md"] = spec_md.read_text(encoding="utf-8")
except OSError as e:
content["spec_md_error"] = str(e)
# Read requirements.json
req = self._load_json(spec_dir / "requirements.json")
if req is not None:
content["requirements"] = req
# Read implementation_plan.json
plan = self._load_json(spec_dir / "implementation_plan.json")
if plan is not None:
content["implementation_plan"] = plan
# Read complexity_assessment.json
assessment = self._load_json(spec_dir / "complexity_assessment.json")
if assessment is not None:
content["complexity_assessment"] = assessment
# Read QA report if present
qa_report = spec_dir / "qa_report.md"
if qa_report.exists():
try:
content["qa_report"] = qa_report.read_text(encoding="utf-8")
except OSError:
pass
return content
def list_specs(self) -> list[dict]:
"""List all specs in the project.
Returns:
List of spec summary dicts
"""
specs_dir = self.project_dir / ".auto-claude" / "specs"
if not specs_dir.is_dir():
return []
specs = []
for item in sorted(specs_dir.iterdir()):
if item.is_dir() and not item.name.startswith("."):
status_info = self.get_spec_status(item.name)
specs.append(status_info)
return specs
def _resolve_spec_dir(self, specs_dir: Path, spec_id: str) -> Path | None:
"""Resolve a spec_id to its directory, supporting prefix matching.
Args:
specs_dir: Parent specs directory
spec_id: Full or prefix spec identifier
Returns:
Path to spec directory or None
"""
# Direct match
exact = specs_dir / spec_id
if exact.is_dir():
return exact
# Prefix match (e.g. '001' matches '001-my-feature')
if specs_dir.is_dir():
for item in specs_dir.iterdir():
if item.is_dir() and item.name.startswith(spec_id):
return item
return None
def _load_json(self, path: Path) -> dict | None:
"""Safely load a JSON file.
Args:
path: Path to the JSON file
Returns:
Parsed dict or None
"""
if not path.exists():
return None
try:
with open(path, encoding="utf-8") as f:
return json.load(f)
except (json.JSONDecodeError, OSError) as e:
logger.warning("Failed to load %s: %s", path, e)
return None
@@ -0,0 +1,429 @@
"""
Task Service
=============
Loads, creates, updates, and deletes tasks by scanning spec directories.
Ported from the TypeScript ProjectStore.loadTasksFromSpecsDir() logic.
"""
from __future__ import annotations
import json
import logging
import re
import shutil
from datetime import datetime, timezone
from pathlib import Path
logger = logging.getLogger(__name__)
# Valid task statuses used by the backend pipeline
VALID_STATUSES = frozenset(
{
"pending",
"spec_creating",
"planning",
"in_progress",
"qa_review",
"qa_fixing",
"human_review",
"done",
"failed",
"cancelled",
}
)
# Status priority for deduplication (higher = more "complete")
_STATUS_PRIORITY: dict[str, int] = {
"done": 100,
"human_review": 80,
"qa_fixing": 70,
"qa_review": 65,
"in_progress": 50,
"planning": 40,
"spec_creating": 35,
"pending": 20,
"cancelled": 15,
"failed": 10,
}
def _slugify(text: str) -> str:
"""Convert a title into a filesystem-safe slug."""
slug = text.lower().strip()
slug = re.sub(r"[^\w\s-]", "", slug)
slug = re.sub(r"[\s_]+", "-", slug)
slug = re.sub(r"-+", "-", slug)
return slug.strip("-")[:80]
def _safe_read_json(path: Path) -> dict | None:
"""Read a JSON file, returning None on any error."""
try:
return json.loads(path.read_text(encoding="utf-8"))
except (json.JSONDecodeError, OSError, ValueError):
return None
def _extract_spec_heading(spec_path: Path) -> str | None:
"""Extract the first markdown heading from a spec.md file."""
try:
content = spec_path.read_text(encoding="utf-8")
match = re.search(
r"^#\s+(?:Quick Spec:|Specification:)?\s*(.+)$", content, re.MULTILINE
)
if match:
return match.group(1).strip()
except OSError:
pass
return None
def _extract_spec_overview(spec_path: Path) -> str | None:
"""Extract the Overview section from a spec.md file."""
try:
content = spec_path.read_text(encoding="utf-8")
match = re.search(r"## Overview\s*\n+([\s\S]*?)(?=\n#{1,6}\s|$)", content)
if match:
return match.group(1).strip()
except OSError:
pass
return None
class TaskService:
"""Manages task lifecycle by reading/writing spec directories."""
def __init__(self, project_dir: Path) -> None:
self.project_dir = project_dir
self.specs_dir = project_dir / ".auto-claude" / "specs"
self.worktrees_dir = project_dir / ".auto-claude" / "worktrees" / "tasks"
# ------------------------------------------------------------------
# Read operations
# ------------------------------------------------------------------
def list_tasks(self) -> list[dict]:
"""Scan spec directories and build a deduplicated task list.
Scans both the main project specs dir and worktree specs dirs.
Main project tasks take priority over worktree duplicates.
"""
all_tasks: list[dict] = []
main_spec_ids: set[str] = set()
# 1. Scan main project specs
if self.specs_dir.is_dir():
main_tasks = self._load_tasks_from_specs_dir(self.specs_dir, "main")
all_tasks.extend(main_tasks)
main_spec_ids = {t["spec_id"] for t in main_tasks}
# 2. Scan worktree specs (only include if spec exists in main)
if self.worktrees_dir.is_dir():
try:
for worktree_dir in sorted(self.worktrees_dir.iterdir()):
if not worktree_dir.is_dir():
continue
wt_specs = worktree_dir / ".auto-claude" / "specs"
if wt_specs.is_dir():
wt_tasks = self._load_tasks_from_specs_dir(wt_specs, "worktree")
valid = [t for t in wt_tasks if t["spec_id"] in main_spec_ids]
all_tasks.extend(valid)
except OSError as exc:
logger.warning("Error scanning worktrees: %s", exc)
# 3. Deduplicate — prefer main over worktree
task_map: dict[str, dict] = {}
for task in all_tasks:
existing = task_map.get(task["spec_id"])
if existing is None:
task_map[task["spec_id"]] = task
else:
existing_is_main = existing.get("location") == "main"
new_is_main = task.get("location") == "main"
if existing_is_main and not new_is_main:
# Keep existing main
continue
elif not existing_is_main and new_is_main:
# Replace worktree with main
task_map[task["spec_id"]] = task
else:
# Same location — use status priority
ep = _STATUS_PRIORITY.get(existing.get("status", ""), 0)
np = _STATUS_PRIORITY.get(task.get("status", ""), 0)
if np > ep:
task_map[task["spec_id"]] = task
return list(task_map.values())
def get_task(self, spec_id: str) -> dict | None:
"""Get full details for a single task by spec_id."""
spec_dir = self.specs_dir / spec_id
if not spec_dir.is_dir():
# Try worktrees
spec_dir = self._find_spec_dir_in_worktrees(spec_id)
if spec_dir is None:
return None
return self._load_single_task(spec_dir, "main")
def create_task(self, title: str, description: str) -> dict:
"""Create a new spec directory with initial files.
Returns the created task dict.
"""
self.specs_dir.mkdir(parents=True, exist_ok=True)
next_num = self._next_spec_number()
slug = _slugify(title)
dir_name = f"{next_num:03d}-{slug}" if slug else f"{next_num:03d}"
spec_dir = self.specs_dir / dir_name
spec_dir.mkdir(parents=True, exist_ok=True)
now = datetime.now(timezone.utc).isoformat()
# Write requirements.json
requirements = {"task_description": description}
(spec_dir / "requirements.json").write_text(
json.dumps(requirements, indent=2), encoding="utf-8"
)
# Write implementation_plan.json
plan = {
"feature": title,
"title": title,
"description": description,
"status": "pending",
"phases": [],
"created_at": now,
"updated_at": now,
}
(spec_dir / "implementation_plan.json").write_text(
json.dumps(plan, indent=2), encoding="utf-8"
)
# Write task_metadata.json
metadata = {
"created_at": now,
"source": "mcp",
}
(spec_dir / "task_metadata.json").write_text(
json.dumps(metadata, indent=2), encoding="utf-8"
)
return self._load_single_task(spec_dir, "main") or {
"spec_id": dir_name,
"title": title,
"description": description,
"status": "pending",
}
def update_task(
self,
spec_id: str,
*,
title: str | None = None,
description: str | None = None,
status: str | None = None,
) -> dict | None:
"""Update task metadata/plan fields."""
spec_dir = self.specs_dir / spec_id
if not spec_dir.is_dir():
return None
plan_path = spec_dir / "implementation_plan.json"
plan = _safe_read_json(plan_path) or {}
changed = False
if title is not None:
plan["feature"] = title
plan["title"] = title
changed = True
if description is not None:
plan["description"] = description
# Also update requirements
req_path = spec_dir / "requirements.json"
reqs = _safe_read_json(req_path) or {}
reqs["task_description"] = description
req_path.write_text(json.dumps(reqs, indent=2), encoding="utf-8")
changed = True
if status is not None:
if status not in VALID_STATUSES:
return None
plan["status"] = status
changed = True
if changed:
plan["updated_at"] = datetime.now(timezone.utc).isoformat()
plan_path.write_text(json.dumps(plan, indent=2), encoding="utf-8")
return self._load_single_task(spec_dir, "main")
def delete_task(self, spec_id: str) -> bool:
"""Delete a spec directory. Returns True if deleted."""
spec_dir = self.specs_dir / spec_id
if not spec_dir.is_dir():
return False
# Safety: ensure it's actually within specs_dir (prevent traversal)
try:
spec_dir.resolve().relative_to(self.specs_dir.resolve())
except ValueError:
logger.error("Path traversal detected for spec_id: %s", spec_id)
return False
shutil.rmtree(spec_dir)
return True
def update_status(self, spec_id: str, status: str) -> dict | None:
"""Update just the status field in implementation_plan.json."""
if status not in VALID_STATUSES:
return None
return self.update_task(spec_id, status=status)
# ------------------------------------------------------------------
# Internal helpers
# ------------------------------------------------------------------
def _next_spec_number(self) -> int:
"""Find the highest existing spec number and return next."""
max_num = 0
if self.specs_dir.is_dir():
for entry in self.specs_dir.iterdir():
if entry.is_dir():
match = re.match(r"^(\d{3})-", entry.name)
if match:
max_num = max(max_num, int(match.group(1)))
return max_num + 1
def _find_spec_dir_in_worktrees(self, spec_id: str) -> Path | None:
"""Search worktree directories for a spec."""
if not self.worktrees_dir.is_dir():
return None
for wt_dir in self.worktrees_dir.iterdir():
if not wt_dir.is_dir():
continue
candidate = wt_dir / ".auto-claude" / "specs" / spec_id
if candidate.is_dir():
return candidate
return None
def _load_tasks_from_specs_dir(self, specs_dir: Path, location: str) -> list[dict]:
"""Load all tasks from a specs directory."""
tasks: list[dict] = []
try:
entries = sorted(specs_dir.iterdir())
except OSError as exc:
logger.warning("Error reading specs directory %s: %s", specs_dir, exc)
return []
for entry in entries:
if not entry.is_dir() or entry.name == ".gitkeep":
continue
try:
task = self._load_single_task(entry, location)
if task:
tasks.append(task)
except Exception as exc:
logger.warning("Error loading spec %s: %s", entry.name, exc)
return tasks
def _load_single_task(self, spec_dir: Path, location: str) -> dict | None:
"""Load a single task from its spec directory."""
dir_name = spec_dir.name
# Read implementation plan
plan = _safe_read_json(spec_dir / "implementation_plan.json")
# Read requirements
requirements = _safe_read_json(spec_dir / "requirements.json")
# Read metadata
metadata = _safe_read_json(spec_dir / "task_metadata.json")
# Determine title (priority: plan.feature > plan.title > dir name)
title = (plan or {}).get("feature") or (plan or {}).get("title") or dir_name
# If title looks like a spec ID (e.g. "054-some-slug"), try spec.md heading
if re.match(r"^\d{3}-", title):
spec_heading = _extract_spec_heading(spec_dir / "spec.md")
if spec_heading:
title = spec_heading
# Determine description (priority: plan.description > requirements.task_description > spec.md overview)
description = ""
if plan and plan.get("description"):
description = plan["description"]
if not description and requirements and requirements.get("task_description"):
description = requirements["task_description"]
if not description:
overview = _extract_spec_overview(spec_dir / "spec.md")
if overview:
description = overview
# Determine status
status = "pending"
if plan and plan.get("status"):
raw_status = plan["status"]
# Map frontend-style statuses to valid backend statuses
status_map: dict[str, str] = {
"pending": "pending",
"backlog": "pending",
"queue": "pending",
"queued": "pending",
"spec_creating": "spec_creating",
"planning": "planning",
"coding": "in_progress",
"in_progress": "in_progress",
"review": "qa_review",
"ai_review": "qa_review",
"qa_review": "qa_review",
"qa_fixing": "qa_fixing",
"human_review": "human_review",
"completed": "done",
"done": "done",
"pr_created": "done",
"error": "failed",
"failed": "failed",
"cancelled": "cancelled",
}
status = status_map.get(raw_status, "pending")
# Extract subtasks from plan phases
subtasks: list[dict] = []
if plan and plan.get("phases"):
for phase in plan["phases"]:
items = phase.get("subtasks") or phase.get("chunks") or []
for st in items:
subtasks.append(
{
"id": st.get("id", ""),
"title": st.get("description", ""),
"status": st.get("status", "pending"),
}
)
# Build result
created_at = (plan or {}).get("created_at", "")
updated_at = (plan or {}).get("updated_at", "")
return {
"spec_id": dir_name,
"title": title,
"description": description,
"status": status,
"subtasks": subtasks,
"metadata": metadata,
"location": location,
"specs_path": str(spec_dir),
"has_spec": (spec_dir / "spec.md").exists(),
"has_plan": (spec_dir / "implementation_plan.json").exists(),
"has_qa_report": (spec_dir / "qa_report.md").exists(),
"created_at": created_at,
"updated_at": updated_at,
}
@@ -0,0 +1,245 @@
"""
Workspace Service
==================
Service layer wrapping the backend WorktreeManager for MCP tool consumption.
Handles git worktree operations: list, diff, merge, discard, and PR creation.
"""
from __future__ import annotations
import contextlib
import io
import logging
from pathlib import Path
logger = logging.getLogger(__name__)
class WorkspaceService:
"""Wraps WorktreeManager operations for MCP server use."""
def __init__(self, project_dir: Path):
self.project_dir = project_dir
def _get_manager(self):
"""Lazily create a WorktreeManager instance.
Returns:
WorktreeManager instance
Raises:
ImportError: If backend module is not available
"""
from core.worktree import WorktreeManager
return WorktreeManager(self.project_dir)
def list_worktrees(self) -> dict:
"""List all active git worktrees for the project.
Returns:
Dict with list of worktree info dicts
"""
try:
manager = self._get_manager()
except ImportError as e:
return {"error": f"Backend module not available: {e}"}
try:
captured = io.StringIO()
with contextlib.redirect_stdout(captured):
worktrees = manager.list_all_worktrees()
result = []
for wt in worktrees:
entry = {
"spec_name": wt.spec_name,
"branch": wt.branch,
"path": str(wt.path),
"base_branch": wt.base_branch,
"is_active": wt.is_active,
"commit_count": wt.commit_count,
"files_changed": wt.files_changed,
"additions": wt.additions,
"deletions": wt.deletions,
}
if wt.days_since_last_commit is not None:
entry["days_since_last_commit"] = wt.days_since_last_commit
if wt.last_commit_date is not None:
entry["last_commit_date"] = wt.last_commit_date.isoformat()
result.append(entry)
return {"worktrees": result, "count": len(result)}
except Exception as e:
logger.exception("Failed to list worktrees")
return {"error": str(e)}
def get_diff(self, spec_id: str) -> dict:
"""Get the git diff for a spec's worktree.
Args:
spec_id: The spec folder name
Returns:
Dict with diff content and change summary
"""
try:
manager = self._get_manager()
except ImportError as e:
return {"error": f"Backend module not available: {e}"}
try:
captured = io.StringIO()
with contextlib.redirect_stdout(captured):
info = manager.get_worktree_info(spec_id)
if info is None:
return {"error": f"No worktree found for spec '{spec_id}'"}
# Get changed files
files = manager.get_changed_files(spec_id)
summary = manager.get_change_summary(spec_id)
# Get actual diff content
from core.git_executable import run_git
diff_result = run_git(
["diff", f"{info.base_branch}...HEAD"],
cwd=info.path,
)
diff_content = ""
if diff_result.returncode == 0:
diff_content = diff_result.stdout
# Truncate very large diffs
if len(diff_content) > 50000:
diff_content = (
diff_content[:50000]
+ "\n\n... (diff truncated, total length: "
+ str(len(diff_result.stdout))
+ " chars)"
)
return {
"spec_id": spec_id,
"branch": info.branch,
"base_branch": info.base_branch,
"changed_files": [
{"status": status, "path": path} for status, path in files
],
"summary": summary,
"diff": diff_content,
}
except Exception as e:
logger.exception("Failed to get diff for %s", spec_id)
return {"error": str(e)}
async def merge(self, spec_id: str, strategy: str = "auto") -> dict:
"""Merge a spec's worktree changes back to the main branch.
Args:
spec_id: The spec folder name
strategy: Merge strategy - 'auto' (git merge), 'no-commit' (stage only)
Returns:
Dict with merge result
"""
try:
manager = self._get_manager()
except ImportError as e:
return {"error": f"Backend module not available: {e}"}
try:
no_commit = strategy == "no-commit"
captured = io.StringIO()
with contextlib.redirect_stdout(captured):
success = manager.merge_worktree(
spec_id,
delete_after=False,
no_commit=no_commit,
)
return {
"success": success,
"spec_id": spec_id,
"strategy": strategy,
"output": captured.getvalue()[-2000:] if captured.getvalue() else "",
}
except Exception as e:
logger.exception("Failed to merge worktree for %s", spec_id)
return {"success": False, "error": str(e)}
def discard(self, spec_id: str) -> dict:
"""Discard a spec's worktree and optionally its branch.
Args:
spec_id: The spec folder name
Returns:
Dict with discard result
"""
try:
manager = self._get_manager()
except ImportError as e:
return {"error": f"Backend module not available: {e}"}
try:
captured = io.StringIO()
with contextlib.redirect_stdout(captured):
manager.remove_worktree(spec_id, delete_branch=True)
return {
"success": True,
"spec_id": spec_id,
"message": f"Worktree and branch for '{spec_id}' removed",
"output": captured.getvalue() if captured.getvalue() else "",
}
except Exception as e:
logger.exception("Failed to discard worktree for %s", spec_id)
return {"success": False, "error": str(e)}
async def create_pr(
self,
spec_id: str,
title: str | None = None,
body: str | None = None,
) -> dict:
"""Push branch and create a pull request from a spec's worktree.
Automatically detects the git provider (GitHub/GitLab).
Args:
spec_id: The spec folder name
title: PR title (defaults to spec name)
body: PR body (defaults to spec summary)
Returns:
Dict with PR URL and status
"""
try:
manager = self._get_manager()
except ImportError as e:
return {"error": f"Backend module not available: {e}"}
try:
captured = io.StringIO()
with contextlib.redirect_stdout(captured):
result = manager.push_and_create_pr(
spec_name=spec_id,
title=title,
)
return {
"success": result.get("success", False),
"spec_id": spec_id,
"pr_url": result.get("pr_url"),
"branch": result.get("branch"),
"provider": result.get("provider"),
"already_exists": result.get("already_exists", False),
"error": result.get("error"),
"output": captured.getvalue()[-1000:] if captured.getvalue() else "",
}
except Exception as e:
logger.exception("Failed to create PR for %s", spec_id)
return {"success": False, "error": str(e)}
@@ -0,0 +1 @@
"""MCP tool modules - each module registers tools with the FastMCP server."""
+159
View File
@@ -0,0 +1,159 @@
"""
Execution Tools
================
MCP tools for starting, stopping, and monitoring builds.
"""
from __future__ import annotations
import asyncio
import logging
from mcp_server.config import get_project_dir
from mcp_server.operations import OperationStatus, tracker
from mcp_server.server import mcp
logger = logging.getLogger(__name__)
@mcp.tool()
async def build_start(
spec_id: str,
model: str = "sonnet",
thinking_level: str = "medium",
) -> dict:
"""Start building/implementing a spec. This is a long-running operation.
Spawns the autonomous coding pipeline which creates a worktree, runs
the planner, then executes each subtask with parallel agents.
Args:
spec_id: Spec folder name or prefix (e.g. '001' or '001-my-feature')
model: Model to use - 'sonnet' (fast), 'opus' (thorough)
thinking_level: Reasoning depth - 'low', 'medium', 'high'
Returns:
An operation_id to poll with operation_get_status() for progress
"""
op = tracker.create("build", f"Building spec: {spec_id}")
async def _run() -> None:
try:
from mcp_server.services.execution_service import ExecutionService
service = ExecutionService(get_project_dir())
tracker.update(
op.id,
status=OperationStatus.RUNNING,
progress=5,
message="Spawning build process...",
)
proc = await service.start_build(
spec_id=spec_id,
model=model,
thinking_level=thinking_level,
)
tracker.update(
op.id,
status=OperationStatus.RUNNING,
progress=10,
message="Build process started, waiting for completion...",
)
# Wait for the process to complete
await proc.wait()
if proc.returncode == 0:
progress_info = service.get_progress(spec_id)
tracker.update(
op.id,
status=OperationStatus.COMPLETED,
progress=100,
message="Build completed successfully",
result=progress_info,
)
else:
logs = service.get_logs(spec_id, tail=20)
tracker.update(
op.id,
status=OperationStatus.FAILED,
error=f"Build exited with code {proc.returncode}",
result={
"exit_code": proc.returncode,
"tail_logs": logs.get("lines", []),
},
)
except Exception as e:
logger.exception("build_start operation failed")
tracker.update(
op.id,
status=OperationStatus.FAILED,
error=str(e),
)
op._task = asyncio.create_task(_run())
return {
"operation_id": op.id,
"message": "Build started. Poll operation_get_status() for progress.",
}
@mcp.tool()
def build_stop(spec_id: str) -> dict:
"""Stop a running build.
Terminates the build subprocess. The worktree and any partial changes
are preserved so the build can be resumed later.
Args:
spec_id: Spec folder name or prefix
Returns:
Whether the build was successfully stopped
"""
from mcp_server.services.execution_service import ExecutionService
service = ExecutionService(get_project_dir())
return service.stop_build(spec_id)
@mcp.tool()
def build_get_progress(spec_id: str) -> dict:
"""Get progress of a running or completed build.
Shows subtask completion status, QA state, and whether the build
is still running.
Args:
spec_id: Spec folder name or prefix
Returns:
Build progress including subtask completion and QA status
"""
from mcp_server.services.execution_service import ExecutionService
service = ExecutionService(get_project_dir())
return service.get_progress(spec_id)
@mcp.tool()
def build_get_logs(spec_id: str, tail: int = 50) -> dict:
"""Get recent build logs for a spec.
Returns the most recent log lines from the build process output.
Args:
spec_id: Spec folder name or prefix
tail: Number of recent lines to return (default 50)
Returns:
Recent log lines from the build
"""
from mcp_server.services.execution_service import ExecutionService
service = ExecutionService(get_project_dir())
return service.get_logs(spec_id, tail=tail)
+208
View File
@@ -0,0 +1,208 @@
"""
GitHub Tools
=============
MCP tools for GitHub automation: PR review, issue triage, auto-fix.
Long-running operations return an operation_id for polling.
"""
from __future__ import annotations
import asyncio
from mcp_server.config import get_project_dir
from mcp_server.operations import OperationStatus, tracker
from mcp_server.server import mcp
from mcp_server.services.github_service import GitHubService
def _get_service() -> GitHubService:
return GitHubService(get_project_dir())
@mcp.tool()
async def github_review_pr(
pr_number: int, repo: str | None = None, model: str = "sonnet"
) -> dict:
"""Review a pull request with AI. Long-running operation - returns operation_id.
Performs a multi-pass AI code review including security, quality,
structural analysis, and AI comment triage.
Args:
pr_number: The PR number to review
repo: Repository in owner/repo format (auto-detected from git remote if omitted)
model: Model to use (haiku, sonnet, opus)
Returns:
Operation ID to poll with operation_get_status()
"""
service = _get_service()
op = tracker.create("github_review_pr", f"Starting review of PR #{pr_number}...")
async def _run():
try:
tracker.update(
op.id,
status=OperationStatus.RUNNING,
progress=10,
message=f"Reviewing PR #{pr_number}...",
)
result = await service.review_pr(pr_number, repo, model)
if "error" in result:
tracker.update(
op.id, status=OperationStatus.FAILED, error=result["error"]
)
else:
tracker.update(
op.id,
status=OperationStatus.COMPLETED,
progress=100,
message="Review complete",
result=result,
)
except Exception as e:
tracker.update(op.id, status=OperationStatus.FAILED, error=str(e))
op._task = asyncio.create_task(_run())
return {
"operation_id": op.id,
"message": f"PR #{pr_number} review started. Poll operation_get_status() for progress.",
}
@mcp.tool()
async def github_list_issues(
state: str = "open", limit: int = 30, repo: str | None = None
) -> dict:
"""List GitHub issues for the project.
Args:
state: Issue state filter: open, closed, or all
limit: Maximum number of issues to return
repo: Repository in owner/repo format (auto-detected if omitted)
Returns:
List of issues with number, title, state, labels, author
"""
service = _get_service()
return await service.list_issues(state, limit, repo)
@mcp.tool()
async def github_auto_fix(issue_number: int, repo: str | None = None) -> dict:
"""Automatically fix a GitHub issue by creating a spec and building it. Long-running.
Creates a specification from the issue, builds it through the autonomous
pipeline (planner -> coder -> QA), and optionally creates a PR.
Args:
issue_number: The issue number to auto-fix
repo: Repository in owner/repo format (auto-detected if omitted)
Returns:
Operation ID to poll with operation_get_status()
"""
service = _get_service()
op = tracker.create(
"github_auto_fix", f"Starting auto-fix for issue #{issue_number}..."
)
async def _run():
try:
tracker.update(
op.id,
status=OperationStatus.RUNNING,
progress=10,
message=f"Auto-fixing issue #{issue_number}...",
)
result = await service.auto_fix_issue(issue_number, repo)
if "error" in result:
tracker.update(
op.id, status=OperationStatus.FAILED, error=result["error"]
)
else:
tracker.update(
op.id,
status=OperationStatus.COMPLETED,
progress=100,
message="Auto-fix complete",
result=result,
)
except Exception as e:
tracker.update(op.id, status=OperationStatus.FAILED, error=str(e))
op._task = asyncio.create_task(_run())
return {
"operation_id": op.id,
"message": f"Auto-fix for issue #{issue_number} started. Poll operation_get_status() for progress.",
}
@mcp.tool()
def github_get_review(pr_number: int) -> dict:
"""Get the most recent review result for a PR.
Returns the saved review data including findings, verdict, and summary.
Args:
pr_number: The PR number to get the review for
Returns:
Review result with findings, verdict, blockers, and summary
"""
service = _get_service()
return service.get_review(pr_number)
@mcp.tool()
async def github_triage_issues(
issue_numbers: list[int], repo: str | None = None
) -> dict:
"""Triage and classify GitHub issues. Long-running.
Analyzes issues for duplicates, spam, feature creep, and assigns
categories, priority, and suggested labels.
Args:
issue_numbers: List of issue numbers to triage
repo: Repository in owner/repo format (auto-detected if omitted)
Returns:
Operation ID to poll with operation_get_status()
"""
service = _get_service()
op = tracker.create(
"github_triage_issues",
f"Starting triage of {len(issue_numbers)} issues...",
)
async def _run():
try:
tracker.update(
op.id,
status=OperationStatus.RUNNING,
progress=10,
message=f"Triaging {len(issue_numbers)} issues...",
)
result = await service.triage_issues(issue_numbers, repo)
if "error" in result:
tracker.update(
op.id, status=OperationStatus.FAILED, error=result["error"]
)
else:
tracker.update(
op.id,
status=OperationStatus.COMPLETED,
progress=100,
message=f"Triaged {result.get('count', 0)} issues",
result=result,
)
except Exception as e:
tracker.update(op.id, status=OperationStatus.FAILED, error=str(e))
op._task = asyncio.create_task(_run())
return {
"operation_id": op.id,
"message": f"Triage of {len(issue_numbers)} issues started. Poll operation_get_status() for progress.",
}
+87
View File
@@ -0,0 +1,87 @@
"""
Ideation Tools
===============
MCP tools for AI-powered project ideation and improvement discovery.
"""
from __future__ import annotations
import asyncio
from mcp_server.config import get_project_dir
from mcp_server.operations import OperationStatus, tracker
from mcp_server.server import mcp
from mcp_server.services.ideation_service import IdeationService
def _get_service() -> IdeationService:
return IdeationService(get_project_dir())
@mcp.tool()
async def ideation_generate(
types: list[str] | None = None,
refresh: bool = False,
model: str = "sonnet",
) -> dict:
"""Generate ideas for project improvements. Long-running.
Analyzes the codebase and generates actionable improvement ideas
across multiple categories.
Args:
types: Ideation types to generate. Options: low_hanging_fruit,
ui_ux_improvements, high_value_features. Defaults to all.
refresh: Force regeneration of existing ideation data
model: Model to use (haiku, sonnet, opus)
Returns:
Operation ID to poll with operation_get_status()
"""
service = _get_service()
op = tracker.create("ideation_generate", "Starting ideation generation...")
async def _run():
try:
tracker.update(
op.id,
status=OperationStatus.RUNNING,
progress=10,
message="Analyzing project for improvement ideas...",
)
result = await service.generate(types=types, refresh=refresh, model=model)
if "error" in result:
tracker.update(
op.id, status=OperationStatus.FAILED, error=result["error"]
)
else:
tracker.update(
op.id,
status=OperationStatus.COMPLETED,
progress=100,
message="Ideation complete",
result=result,
)
except Exception as e:
tracker.update(op.id, status=OperationStatus.FAILED, error=str(e))
op._task = asyncio.create_task(_run())
return {
"operation_id": op.id,
"message": "Ideation generation started. Poll operation_get_status() for progress.",
}
@mcp.tool()
def ideation_get() -> dict:
"""Get previously generated ideation results.
Returns all generated ideas with their categories, priorities,
effort estimates, and implementation suggestions.
Returns:
Ideation data with ideas grouped by type and priority
"""
service = _get_service()
return service.get_ideation()
+84
View File
@@ -0,0 +1,84 @@
"""
Insights Tools
===============
MCP tools for AI-powered codebase insights and Q&A.
"""
from __future__ import annotations
import asyncio
from mcp_server.config import get_project_dir
from mcp_server.operations import OperationStatus, tracker
from mcp_server.server import mcp
from mcp_server.services.insights_service import InsightsService
def _get_service() -> InsightsService:
return InsightsService(get_project_dir())
@mcp.tool()
async def insights_ask(
question: str, history: list | None = None, model: str = "sonnet"
) -> dict:
"""Ask an AI question about the codebase. Long-running operation.
The AI agent has access to the codebase and can read files, search,
and explore to answer questions about architecture, patterns, bugs, etc.
Args:
question: The question to ask about the codebase
history: Optional conversation history as list of {role, content} dicts
model: Model to use (haiku, sonnet, opus)
Returns:
Operation ID to poll with operation_get_status()
"""
service = _get_service()
op = tracker.create("insights_ask", "Processing question...")
async def _run():
try:
tracker.update(
op.id,
status=OperationStatus.RUNNING,
progress=10,
message="AI is exploring the codebase...",
)
result = await service.ask(question, history, model)
if "error" in result:
tracker.update(
op.id, status=OperationStatus.FAILED, error=result["error"]
)
else:
tracker.update(
op.id,
status=OperationStatus.COMPLETED,
progress=100,
message="Question answered",
result=result,
)
except Exception as e:
tracker.update(op.id, status=OperationStatus.FAILED, error=str(e))
op._task = asyncio.create_task(_run())
return {
"operation_id": op.id,
"message": "Insights query started. Poll operation_get_status() for the answer.",
}
@mcp.tool()
def insights_suggest_tasks() -> dict:
"""Get AI-suggested tasks based on recent insights conversations.
Returns task suggestions derived from ideation data or previous
insights conversations.
Returns:
List of task suggestions with title, description, category, impact
"""
service = _get_service()
return service.suggest_tasks()
+75
View File
@@ -0,0 +1,75 @@
"""
Memory Tools
=============
MCP tools for Graphiti-based semantic memory (knowledge graph).
Requires GRAPHITI_ENABLED=true in the environment.
"""
from __future__ import annotations
from mcp_server.config import get_project_dir
from mcp_server.server import mcp
from mcp_server.services.memory_service import MemoryService
def _get_service() -> MemoryService:
return MemoryService(get_project_dir())
@mcp.tool()
async def memory_search(query: str, limit: int = 10) -> dict:
"""Search the project's semantic memory (Graphiti knowledge graph).
Finds relevant stored knowledge including codebase discoveries,
session insights, patterns, gotchas, and task outcomes.
Requires GRAPHITI_ENABLED=true in environment.
Args:
query: Search query describing what you're looking for
limit: Maximum number of results to return (default 10)
Returns:
List of relevant memory entries with content and relevance scores
"""
service = _get_service()
return await service.search(query, limit)
@mcp.tool()
async def memory_add_episode(content: str, source: str = "mcp") -> dict:
"""Add a new episode/fact to the project's memory.
Stores information in the knowledge graph for future retrieval.
Use this to record insights, patterns, or important findings.
Requires GRAPHITI_ENABLED=true in environment.
Args:
content: The information to store (insight, pattern, discovery, etc.)
source: Source identifier for the episode (default: mcp)
Returns:
Confirmation of successful storage
"""
service = _get_service()
return await service.add_episode(content, source)
@mcp.tool()
async def memory_get_recent(limit: int = 10) -> dict:
"""Get recent memory entries.
Retrieves the most recent entries from the project's knowledge graph.
Requires GRAPHITI_ENABLED=true in environment.
Args:
limit: Maximum number of entries to return (default 10)
Returns:
List of recent memory entries
"""
service = _get_service()
return await service.get_recent(limit)
+52
View File
@@ -0,0 +1,52 @@
"""
Operations Management Tools
============================
Tools for polling long-running operation status and cancelling operations.
"""
from __future__ import annotations
from mcp_server.operations import tracker
from mcp_server.server import mcp
@mcp.tool()
def operation_get_status(operation_id: str) -> dict:
"""Get the status of a long-running operation.
Use this to poll for progress on operations started by tools like
spec_create, build_start, qa_start_review, etc.
Args:
operation_id: The operation ID returned by the tool that started the operation
Returns:
Operation status including progress (0-100), message, and result when complete
"""
op = tracker.get(operation_id)
if op is None:
return {"error": f"Operation {operation_id} not found"}
return op.to_dict()
@mcp.tool()
def operation_cancel(operation_id: str) -> dict:
"""Cancel a running operation.
Args:
operation_id: The operation ID to cancel
Returns:
Whether the cancellation was successful
"""
success = tracker.cancel(operation_id)
if not success:
op = tracker.get(operation_id)
if op is None:
return {"success": False, "error": "Operation not found"}
return {
"success": False,
"error": f"Cannot cancel operation in {op.status.value} state",
}
return {"success": True, "message": "Operation cancelled"}
+148
View File
@@ -0,0 +1,148 @@
"""
Project Management Tools
=========================
MCP tools for managing the active project: switching projects,
getting status, listing specs, and reading the project index.
"""
from __future__ import annotations
import logging
from mcp_server import config
from mcp_server.server import mcp
logger = logging.getLogger(__name__)
@mcp.tool()
def project_set_active(project_dir: str) -> dict:
"""Switch the MCP server to a different project directory.
Re-initializes the server to point at a new project.
All subsequent tool calls will operate on this project.
Args:
project_dir: Absolute path to the project directory
"""
try:
config.initialize(project_dir)
project_path = config.get_project_dir()
initialized = config.is_initialized()
return {
"success": True,
"project_dir": str(project_path),
"initialized": initialized,
"message": f"Active project set to {project_path}",
}
except (ValueError, RuntimeError) as exc:
return {"success": False, "error": str(exc)}
@mcp.tool()
def project_get_status() -> dict:
"""Get the current project status.
Returns the active project directory, initialization state,
specs count, and project index summary.
"""
try:
project_dir = config.get_project_dir()
except RuntimeError:
return {
"initialized": False,
"error": "No project set. Use project_set_active() first.",
}
initialized = config.is_initialized()
specs_count = 0
if initialized:
specs_dir = config.get_specs_dir()
if specs_dir.is_dir():
specs_count = sum(
1
for entry in specs_dir.iterdir()
if entry.is_dir() and entry.name != ".gitkeep"
)
index = config.get_project_index()
return {
"project_dir": str(project_dir),
"initialized": initialized,
"specs_count": specs_count,
"has_project_index": bool(index),
"project_name": index.get("name", project_dir.name),
}
@mcp.tool()
def project_list_specs() -> dict:
"""List all spec directories with their basic info.
Returns each spec's name, whether it has a plan/spec file,
and status from the implementation plan.
"""
try:
specs_dir = config.get_specs_dir()
except RuntimeError:
return {"error": "No project set. Use project_set_active() first.", "specs": []}
if not specs_dir.is_dir():
return {"specs": [], "message": "No specs directory found."}
specs: list[dict] = []
for entry in sorted(specs_dir.iterdir()):
if not entry.is_dir() or entry.name == ".gitkeep":
continue
has_plan = (entry / "implementation_plan.json").exists()
has_spec = (entry / "spec.md").exists()
status = "pending"
title = entry.name
if has_plan:
try:
import json
plan = json.loads(
(entry / "implementation_plan.json").read_text(encoding="utf-8")
)
status = plan.get("status", "pending")
title = plan.get("feature") or plan.get("title") or entry.name
except (json.JSONDecodeError, OSError):
pass
specs.append(
{
"name": entry.name,
"title": title,
"has_plan": has_plan,
"has_spec": has_spec,
"has_qa_report": (entry / "qa_report.md").exists(),
"status": status,
}
)
return {"specs": specs, "count": len(specs)}
@mcp.tool()
def project_get_index() -> dict:
"""Return the full project_index.json content.
The project index contains metadata about the project
such as file summaries, dependency info, and analysis results.
"""
try:
index = config.get_project_index()
except RuntimeError:
return {"error": "No project set. Use project_set_active() first."}
if not index:
return {"message": "No project index found. Run indexing first.", "index": {}}
return {"index": index}
+125
View File
@@ -0,0 +1,125 @@
"""
QA Tools
=========
MCP tools for running QA reviews, getting reports, and manual approval.
"""
from __future__ import annotations
import asyncio
import logging
from mcp_server.config import get_project_dir
from mcp_server.operations import OperationStatus, tracker
from mcp_server.server import mcp
logger = logging.getLogger(__name__)
@mcp.tool()
async def qa_start_review(spec_id: str) -> dict:
"""Start QA review for a completed build. This is a long-running operation.
Runs the QA reviewer agent which validates the implementation against
the spec's acceptance criteria. The agent reads code, runs tests, and
produces a detailed QA report.
Args:
spec_id: Spec folder name or prefix (e.g. '001' or '001-my-feature')
Returns:
An operation_id to poll with operation_get_status() for progress
"""
op = tracker.create("qa_review", f"QA review for: {spec_id}")
async def _run() -> None:
try:
from mcp_server.services.qa_service import QAService
service = QAService(get_project_dir())
tracker.update(
op.id,
status=OperationStatus.RUNNING,
progress=10,
message="Starting QA review session...",
)
result = await service.start_review(spec_id=spec_id)
status = result.get("status", "error")
if status == "approved":
tracker.update(
op.id,
status=OperationStatus.COMPLETED,
progress=100,
message="QA approved - all acceptance criteria validated",
result=result,
)
elif status == "rejected":
tracker.update(
op.id,
status=OperationStatus.COMPLETED,
progress=100,
message="QA rejected - issues found, see report",
result=result,
)
else:
tracker.update(
op.id,
status=OperationStatus.FAILED,
error=result.get("error", "QA review failed"),
result=result,
)
except Exception as e:
logger.exception("qa_start_review operation failed")
tracker.update(
op.id,
status=OperationStatus.FAILED,
error=str(e),
)
op._task = asyncio.create_task(_run())
return {
"operation_id": op.id,
"message": "QA review started. Poll operation_get_status() for progress.",
}
@mcp.tool()
def qa_get_report(spec_id: str) -> dict:
"""Get the QA report for a spec.
Returns the full QA report including validation results, issues found,
and the qa_signoff status from the implementation plan.
Args:
spec_id: Spec folder name or prefix
Returns:
QA report content, fix requests, and signoff status
"""
from mcp_server.services.qa_service import QAService
service = QAService(get_project_dir())
return service.get_report(spec_id)
@mcp.tool()
def qa_approve(spec_id: str) -> dict:
"""Manually approve a spec that's in QA review.
Use this to bypass the automated QA review and mark a spec as approved.
This updates the implementation_plan.json qa_signoff status.
Args:
spec_id: Spec folder name or prefix
Returns:
Whether the approval was successful
"""
from mcp_server.services.qa_service import QAService
service = QAService(get_project_dir())
return service.approve(spec_id)
+128
View File
@@ -0,0 +1,128 @@
"""
Roadmap Tools
==============
MCP tools for AI-powered strategic roadmap generation.
"""
from __future__ import annotations
import asyncio
from mcp_server.config import get_project_dir
from mcp_server.operations import OperationStatus, tracker
from mcp_server.server import mcp
from mcp_server.services.roadmap_service import RoadmapService
def _get_service() -> RoadmapService:
return RoadmapService(get_project_dir())
@mcp.tool()
async def roadmap_generate(refresh: bool = False, model: str = "sonnet") -> dict:
"""Generate a strategic roadmap for the project. Long-running.
Analyzes the project structure, existing features, and codebase to
generate a phased roadmap with prioritized features.
Args:
refresh: Force regeneration even if a roadmap already exists
model: Model to use (haiku, sonnet, opus)
Returns:
Operation ID to poll with operation_get_status()
"""
service = _get_service()
op = tracker.create("roadmap_generate", "Starting roadmap generation...")
async def _run():
try:
tracker.update(
op.id,
status=OperationStatus.RUNNING,
progress=10,
message="Analyzing project for roadmap generation...",
)
result = await service.generate(refresh=refresh, model=model)
if "error" in result:
tracker.update(
op.id, status=OperationStatus.FAILED, error=result["error"]
)
else:
tracker.update(
op.id,
status=OperationStatus.COMPLETED,
progress=100,
message="Roadmap generated",
result=result,
)
except Exception as e:
tracker.update(op.id, status=OperationStatus.FAILED, error=str(e))
op._task = asyncio.create_task(_run())
return {
"operation_id": op.id,
"message": "Roadmap generation started. Poll operation_get_status() for progress.",
}
@mcp.tool()
def roadmap_get() -> dict:
"""Get the current roadmap data.
Returns the previously generated roadmap including vision, phases,
features with priorities, and implementation details.
Returns:
Roadmap data with vision, phases, features, and priority breakdown
"""
service = _get_service()
return service.get_roadmap()
@mcp.tool()
async def roadmap_refresh(model: str = "sonnet") -> dict:
"""Refresh/regenerate the roadmap. Long-running.
Forces a complete regeneration of the roadmap, analyzing current
project state and creating updated phases and features.
Args:
model: Model to use (haiku, sonnet, opus)
Returns:
Operation ID to poll with operation_get_status()
"""
service = _get_service()
op = tracker.create("roadmap_refresh", "Starting roadmap refresh...")
async def _run():
try:
tracker.update(
op.id,
status=OperationStatus.RUNNING,
progress=10,
message="Refreshing roadmap...",
)
result = await service.generate(refresh=True, model=model)
if "error" in result:
tracker.update(
op.id, status=OperationStatus.FAILED, error=result["error"]
)
else:
tracker.update(
op.id,
status=OperationStatus.COMPLETED,
progress=100,
message="Roadmap refreshed",
result=result,
)
except Exception as e:
tracker.update(op.id, status=OperationStatus.FAILED, error=str(e))
op._task = asyncio.create_task(_run())
return {
"operation_id": op.id,
"message": "Roadmap refresh started. Poll operation_get_status() for progress.",
}
+149
View File
@@ -0,0 +1,149 @@
"""
Spec Tools
===========
MCP tools for creating and inspecting specifications.
"""
from __future__ import annotations
import asyncio
import logging
from mcp_server.config import get_project_dir
from mcp_server.operations import OperationStatus, tracker
from mcp_server.server import mcp
logger = logging.getLogger(__name__)
@mcp.tool()
async def spec_create(
task_description: str,
model: str = "sonnet",
thinking_level: str = "medium",
) -> dict:
"""Create a specification for a task. This is a long-running operation.
The spec creation pipeline runs multiple AI phases (discovery, requirements,
complexity assessment, spec writing, planning) to produce a complete
implementation-ready specification.
Args:
task_description: What you want to build (be specific and detailed)
model: Model to use - 'sonnet' (fast), 'opus' (thorough)
thinking_level: How much the AI reasons - 'low', 'medium', 'high'
Returns:
An operation_id to poll with operation_get_status() for progress
"""
op = tracker.create("spec_create", f"Creating spec for: {task_description[:80]}")
async def _run() -> None:
try:
from mcp_server.services.spec_service import SpecService
tracker.update(
op.id,
status=OperationStatus.RUNNING,
progress=5,
message="Initializing spec pipeline...",
)
service = SpecService(get_project_dir())
tracker.update(
op.id,
status=OperationStatus.RUNNING,
progress=10,
message="Running spec creation phases...",
)
result = await service.create_spec(
task_description=task_description,
model=model,
thinking_level=thinking_level,
)
if result.get("success"):
tracker.update(
op.id,
status=OperationStatus.COMPLETED,
progress=100,
message="Spec created successfully",
result=result,
)
else:
tracker.update(
op.id,
status=OperationStatus.FAILED,
progress=100,
message="Spec creation failed",
error=result.get("error", "Unknown error"),
result=result,
)
except Exception as e:
logger.exception("spec_create operation failed")
tracker.update(
op.id,
status=OperationStatus.FAILED,
error=str(e),
)
op._task = asyncio.create_task(_run())
return {
"operation_id": op.id,
"message": "Spec creation started. Poll operation_get_status() for progress.",
}
@mcp.tool()
def spec_get_status(spec_id: str) -> dict:
"""Get the current status of a spec (which phases have completed).
Shows whether discovery, requirements, spec writing, and planning
phases are complete, and the overall readiness state.
Args:
spec_id: Spec folder name or prefix (e.g. '001' or '001-my-feature')
Returns:
Status including completed phases and overall state
"""
from mcp_server.services.spec_service import SpecService
service = SpecService(get_project_dir())
return service.get_spec_status(spec_id)
@mcp.tool()
def spec_get_content(spec_id: str) -> dict:
"""Get the full content of a spec including spec.md, requirements, and plan.
Returns the complete specification content so you can understand what
will be built and how.
Args:
spec_id: Spec folder name or prefix (e.g. '001' or '001-my-feature')
Returns:
Full spec content including spec.md, requirements.json, implementation_plan.json
"""
from mcp_server.services.spec_service import SpecService
service = SpecService(get_project_dir())
return service.get_spec_content(spec_id)
@mcp.tool()
def spec_list() -> dict:
"""List all specs in the project with their status.
Returns:
List of specs with their current status and phase completion
"""
from mcp_server.services.spec_service import SpecService
service = SpecService(get_project_dir())
specs = service.list_specs()
return {"specs": specs, "count": len(specs)}
+179
View File
@@ -0,0 +1,179 @@
"""
Task Management Tools
======================
MCP tools for CRUD operations on tasks (specs).
Tasks are stored as spec directories under .auto-claude/specs/.
"""
from __future__ import annotations
import logging
from mcp_server import config
from mcp_server.server import mcp
from mcp_server.services.task_service import VALID_STATUSES, TaskService
logger = logging.getLogger(__name__)
def _get_task_service() -> TaskService:
"""Get a TaskService instance for the active project."""
return TaskService(config.get_project_dir())
@mcp.tool()
def task_list() -> dict:
"""List all tasks with id, title, status, and description preview.
Scans both the main project and worktree spec directories,
deduplicating with main project taking priority.
"""
try:
service = _get_task_service()
except RuntimeError as exc:
return {"error": str(exc), "tasks": []}
tasks = service.list_tasks()
# Return a concise view for listing
summary = []
for t in tasks:
desc = t.get("description", "")
preview = (desc[:200] + "...") if len(desc) > 200 else desc
summary.append(
{
"spec_id": t["spec_id"],
"title": t["title"],
"status": t["status"],
"description_preview": preview,
"has_spec": t.get("has_spec", False),
"has_plan": t.get("has_plan", False),
"subtask_count": len(t.get("subtasks", [])),
}
)
return {"tasks": summary, "count": len(summary)}
@mcp.tool()
def task_create(title: str, description: str) -> dict:
"""Create a new task (spec directory) with initial files.
Generates the next spec number automatically and creates
the directory with requirements.json, implementation_plan.json,
and task_metadata.json.
Args:
title: The task title (used for the directory name slug)
description: Full description of what needs to be done
"""
try:
service = _get_task_service()
except RuntimeError as exc:
return {"error": str(exc)}
if not title or not title.strip():
return {"error": "Title is required"}
if not description or not description.strip():
return {"error": "Description is required"}
task = service.create_task(title.strip(), description.strip())
return {"success": True, "task": task}
@mcp.tool()
def task_get(spec_id: str) -> dict:
"""Get full task details including subtasks, metadata, and file info.
Args:
spec_id: The spec directory name (e.g., "001-my-feature")
"""
try:
service = _get_task_service()
except RuntimeError as exc:
return {"error": str(exc)}
task = service.get_task(spec_id)
if task is None:
return {"error": f"Task '{spec_id}' not found"}
return {"task": task}
@mcp.tool()
def task_update(
spec_id: str,
title: str | None = None,
description: str | None = None,
status: str | None = None,
) -> dict:
"""Update task metadata (title, description, and/or status).
Args:
spec_id: The spec directory name (e.g., "001-my-feature")
title: New title (optional)
description: New description (optional)
status: New status (optional) - must be a valid status
"""
try:
service = _get_task_service()
except RuntimeError as exc:
return {"error": str(exc)}
if status is not None and status not in VALID_STATUSES:
return {"error": f"Invalid status '{status}'. Valid: {sorted(VALID_STATUSES)}"}
task = service.update_task(
spec_id, title=title, description=description, status=status
)
if task is None:
return {"error": f"Task '{spec_id}' not found or invalid update"}
return {"success": True, "task": task}
@mcp.tool()
def task_delete(spec_id: str) -> dict:
"""Delete a task by removing its spec directory.
WARNING: This permanently deletes the spec directory and all its contents.
Args:
spec_id: The spec directory name (e.g., "001-my-feature")
"""
try:
service = _get_task_service()
except RuntimeError as exc:
return {"error": str(exc)}
deleted = service.delete_task(spec_id)
if not deleted:
return {"error": f"Task '{spec_id}' not found"}
return {"success": True, "message": f"Task '{spec_id}' deleted"}
@mcp.tool()
def task_update_status(spec_id: str, status: str) -> dict:
"""Update just the status field of a task.
Args:
spec_id: The spec directory name (e.g., "001-my-feature")
status: New status value. Valid statuses: pending, spec_creating,
planning, in_progress, qa_review, qa_fixing, human_review,
done, failed, cancelled
"""
try:
service = _get_task_service()
except RuntimeError as exc:
return {"error": str(exc)}
if status not in VALID_STATUSES:
return {"error": f"Invalid status '{status}'. Valid: {sorted(VALID_STATUSES)}"}
task = service.update_status(spec_id, status)
if task is None:
return {"error": f"Task '{spec_id}' not found"}
return {"success": True, "task": task}
+159
View File
@@ -0,0 +1,159 @@
"""
Workspace Tools
================
MCP tools for managing git worktrees: list, diff, merge, discard, and PR creation.
"""
from __future__ import annotations
import asyncio
import logging
from mcp_server.config import get_project_dir
from mcp_server.operations import OperationStatus, tracker
from mcp_server.server import mcp
logger = logging.getLogger(__name__)
@mcp.tool()
def workspace_list() -> dict:
"""List all active git worktrees for the project.
Each spec gets its own isolated worktree. This shows all active
worktrees with their branch, change stats, and age.
Returns:
List of worktrees with branch, stats, and age information
"""
from mcp_server.services.workspace_service import WorkspaceService
service = WorkspaceService(get_project_dir())
return service.list_worktrees()
@mcp.tool()
def workspace_diff(spec_id: str) -> dict:
"""Get the git diff for a spec's worktree.
Shows all changes made in the spec's branch compared to the base branch,
including a file-level summary and the full diff content.
Args:
spec_id: Spec folder name or prefix (e.g. '001' or '001-my-feature')
Returns:
Changed files summary and full diff content
"""
from mcp_server.services.workspace_service import WorkspaceService
service = WorkspaceService(get_project_dir())
return service.get_diff(spec_id)
@mcp.tool()
async def workspace_merge(spec_id: str, strategy: str = "auto") -> dict:
"""Merge a spec's worktree changes back to the main branch.
Args:
spec_id: Spec folder name or prefix
strategy: 'auto' for standard git merge, 'no-commit' to stage without committing
Returns:
Whether the merge was successful
"""
from mcp_server.services.workspace_service import WorkspaceService
service = WorkspaceService(get_project_dir())
return await service.merge(spec_id, strategy=strategy)
@mcp.tool()
def workspace_discard(spec_id: str) -> dict:
"""Discard a spec's worktree and its branch.
Permanently removes the worktree directory and deletes the associated
git branch. This cannot be undone.
Args:
spec_id: Spec folder name or prefix
Returns:
Whether the discard was successful
"""
from mcp_server.services.workspace_service import WorkspaceService
service = WorkspaceService(get_project_dir())
return service.discard(spec_id)
@mcp.tool()
async def workspace_create_pr(
spec_id: str,
title: str | None = None,
body: str | None = None,
) -> dict:
"""Create a pull request from a spec's worktree branch.
Pushes the branch to origin and creates a PR/MR on the detected
git hosting provider (GitHub or GitLab). This is a long-running operation.
Args:
spec_id: Spec folder name or prefix
title: PR title (defaults to spec name)
body: PR body (defaults to spec summary)
Returns:
An operation_id to poll with operation_get_status() for progress
"""
op = tracker.create("create_pr", f"Creating PR for: {spec_id}")
async def _run() -> None:
try:
from mcp_server.services.workspace_service import WorkspaceService
service = WorkspaceService(get_project_dir())
tracker.update(
op.id,
status=OperationStatus.RUNNING,
progress=20,
message="Pushing branch and creating PR...",
)
result = await service.create_pr(
spec_id=spec_id,
title=title,
body=body,
)
if result.get("success"):
pr_url = result.get("pr_url", "")
tracker.update(
op.id,
status=OperationStatus.COMPLETED,
progress=100,
message=f"PR created: {pr_url}" if pr_url else "PR created",
result=result,
)
else:
tracker.update(
op.id,
status=OperationStatus.FAILED,
error=result.get("error", "PR creation failed"),
result=result,
)
except Exception as e:
logger.exception("workspace_create_pr operation failed")
tracker.update(
op.id,
status=OperationStatus.FAILED,
error=str(e),
)
op._task = asyncio.create_task(_run())
return {
"operation_id": op.id,
"message": "PR creation started. Poll operation_get_status() for progress.",
}