6.3 KiB
CLAUDE.md
This file provides guidance to Claude Code (claude.ai/code) when working with code in this repository.
Project Overview
spark-executor-mcp is a Python 3.12+ service that exposes Spark-on-YARN operations as MCP tools via fastapi-mcp, so an LLM agent can submit, monitor, fetch logs from, and kill PySpark jobs. Submission is a deliberate two-step flow (prepare_submit_job then confirm_submit_job on second confirmation), with a saved Connection registry (so the agent can pick which YARN/standalone cluster to use).
Common Commands
The project uses uv for dependency management. Dependencies are pinned in pyproject.toml and uv.lock. A virtualenv already exists at .venv/. PyPI index is configured to the Tsinghua mirror in pyproject.toml.
# Install/sync dependencies
uv sync
# Run the server (binds 0.0.0.0:8000; MCP endpoint at /spark-executor-mcp)
uv run main.py
# Run the full test suite
uv run pytest
# Run a single test file
uv run pytest tests/unit/test_submit_tool.py -v
# Run tests matching a name
uv run pytest -v -k cancel_pending
# Activate venv and run directly
source .venv/bin/activate
python main.py
The MCP transport needs an MCP initialize handshake first to get a mcp-session-id header; only then do tools/list / tools/call work.
Architecture
main.py # Root FastAPI app; lifespan wires in the MCP server
common/
factory.py # init_mcp_server(app) -> FastApiMCP wrapper
logging.py # loguru: stderr + data/logs/{debug,info}/*.log, rotated 30d
spark_executor/
__init__.py # Re-exports `app`
server.py # FastAPI app; 12 tool routes + exception handlers
models.py # Pydantic: Job, JobStatus, SubmitResult, Connection, PendingSubmission
tools/
submit.py # prepare / confirm / list / get / cancel pending + job_store
status.py logs.py kill.py # job-lifecycle tools
connections.py # save / list / get / delete connection tools
requests.py # Pydantic body models for every FastAPI route
core/
spark_submit.py # builds & runs spark-submit commands
yarn_client.py # wraps yarn application / yarn logs
log_parser.py # extracts application_id from spark-submit output
job_store.py # in-memory dict: job_id -> Job (Stage 3 -> SQLite)
connection_store.py # JSON CRUD over ./data/connections.json
pending_store.py # JSON CRUD over ./data/pending_jobs.json
data/ # gitignored: connections.json, pending_jobs.json, logs/
tests/ # unit/ + integration/; conftest adds repo root to sys.path
docs/superpowers/plans/ # implementation plans
How a request flows:
main.pycreates the rootFastAPI(title="Main App")and registers anasynccontextmanagerlifespan.- At startup, the lifespan calls
init_mcp_server(spark_executor_app)(common/factory.py) which returns aFastApiMCPinstance bound to the spark executor's FastAPI app. - That
FastApiMCPis mounted onto the root app at/spark-executor-mcpviamount_http(...).fastapi-mcpinspects the spark executor's routes and registers each one as an MCP tool. - An MCP client calls
initialize(gets amcp-session-id), thentools/list(sees all 12 tools), thentools/call(sends args as a JSON body, FastAPI validates via the Pydantic model intools/requests.py).
Persistence layout (under ./data/, overridable via SPARK_EXECUTOR_DATA_DIR):
data/connections.json—Connectionrecords keyed by name (atomic write viatempfile+os.replace)data/pending_jobs.json—PendingSubmissionrecords keyed bypending_iddata/logs/debug/YYYY-MM-DD.log— DEBUG sink, gzipped, 30-day retentiondata/logs/info/YYYY-MM-DD.log— INFO sink, gzipped, 30-day retentionJobStoreis in-memory only today; Stage 3 will move it to SQLite.
Key conventions:
- All modules include the file header
# coding=utf-8plus a@Time/@Authordocstring. Match this when adding new files. common/logging.pyconfigures a single process-widelogurulogger (stderr + two rotated file sinks). Importfrom common.logging import loggerrather than creating new loggers.spark_executor/__init__.pyre-exportsappso callers canfrom spark_executor import app— keep this re-export when adding to the package.- The
FastApiMCPinstance is created per-app in the lifespan; do not cache it at module import time. - FastAPI routes use Pydantic body models (from
tools/requests.py), not query parameters.fastapi-mcppasses tool args as a JSON body, and dict-typed query params arrive as strings and 422. Every new MCP tool needs a request model intools/requests.py. - Exception handlers in
server.py:KeyError→ 404,ValueError→ 400, Pydantic validation → 422. Use the existing handlers — don't add try/except in route bodies. - Tool functions raise
KeyErrorfor "unknown id" andValueErrorfor invalid state transitions (e.g. confirming a CANCELLED pending). The handlers translate these to clean HTTP statuses. - The two-step submit flow is core, not optional.
prepare_submit_jobsnapshots the connection'smaster/deploy_mode/spark_confinto thePendingSubmission;confirm_submit_jobis the only placespark-submitis invoked. Editing a connection between prepare and confirm does not retarget the pending job. - Test fixtures rebind module-level singletons (
connections.store,submit.conn_store,submit.pending_store) inmonkeypatch.setattrbecause the tool modules captured the originals at import time. Seetests/integration/test_mcp_routes.pyfor the pattern. - 86 tests pass as of the last Stage 1 cleanup; run
uv run pytestafter any change.
Stage Status
- Stage 1: complete — 12 MCP tools, 2-step submit, connection management, structured logging, Pydantic body models, exception handlers. 22 commits on
feat/stage-1. - Stage 2: not started —
generate_job_filefor LLM-written PySpark code (job_writer+ one new tool). - Stage 3: deferred — async submission, SQLite
JobRegistry, status poller, offset-basedLogCache, multi-tenantowner. Plan lives atdocs/superpowers/plans/2026-06-24-spark-executor-mcp.md.