From ce2eeb4bac99a4a9e62cde0e85c7832948f0c3f4 Mon Sep 17 00:00:00 2001 From: Claude Date: Wed, 24 Jun 2026 15:20:29 +0800 Subject: [PATCH] feat: update memory --- CLAUDE.md | 64 ++++++++++++++++++++++++++++++++++++++++++++----------- 1 file changed, 52 insertions(+), 12 deletions(-) diff --git a/CLAUDE.md b/CLAUDE.md index cb78593..a7cac83 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -4,7 +4,7 @@ This file provides guidance to Claude Code (claude.ai/code) when working with co ## Project Overview -`spark-executor-mcp` is a Python 3.12+ service that wraps Spark-related operations as a **Model Context Protocol (MCP)** server, exposed via HTTP through FastAPI. The MCP layer is built on top of [`fastapi-mcp`](https://github.com/tadata-org/fastapi_mcp), which auto-discovers FastAPI routes and exposes them as MCP tools. +`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 @@ -14,30 +14,51 @@ The project uses `uv` for dependency management. Dependencies are pinned in `pyp # Install/sync dependencies uv sync -# Run the server (binds 0.0.0.0:8000) +# Run the server (binds 0.0.0.0:8000; MCP endpoint at /spark-executor-mcp) uv run main.py -# Or activate the venv and run directly +# 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 ``` -Once running, the MCP endpoint is mounted at `http://localhost:8000/spark-executor-mcp`. The `spark_executor` app's own routes (e.g. `/health`) remain available at the root. - -There is currently no test suite, linter configuration, or CI config in the repository. +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 -The codebase is small and intentionally layered — the goal is to add new "tool" FastAPI sub-apps and have them be auto-exposed as MCP. - ``` main.py # Root FastAPI app; lifespan wires in the MCP server common/ factory.py # init_mcp_server(app) -> FastApiMCP wrapper - logging.py # loguru logger, level=DEBUG to stderr + logging.py # loguru: stderr + data/logs/{debug,info}/*.log, rotated 30d spark_executor/ __init__.py # Re-exports `app` - server.py # Spark Executor FastAPI app (the MCP tool surface) + 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:** @@ -45,11 +66,30 @@ spark_executor/ 1. `main.py` creates the root `FastAPI(title="Main App")` and registers an `asynccontextmanager` `lifespan`. 2. At startup, the lifespan calls `init_mcp_server(spark_executor_app)` (`common/factory.py`) which returns a `FastApiMCP` instance bound to the spark executor's FastAPI app. 3. That `FastApiMCP` is mounted onto the root app at `/spark-executor-mcp` via `mount_http(...)`. `fastapi-mcp` inspects the spark executor's routes and registers each one as an MCP tool. -4. To add new MCP tools, define new routes on the `app` in `spark_executor/server.py` (or create a new sub-package and mount it the same way in `main.py`'s lifespan) — they will be picked up automatically. +4. An MCP client calls `initialize` (gets a `mcp-session-id`), then `tools/list` (sees all 12 tools), then `tools/call` (sends args as a JSON body, FastAPI validates via the Pydantic model in `tools/requests.py`). + +**Persistence layout** (under `./data/`, overridable via `SPARK_EXECUTOR_DATA_DIR`): +- `data/connections.json` — `Connection` records keyed by name (atomic write via `tempfile` + `os.replace`) +- `data/pending_jobs.json` — `PendingSubmission` records keyed by `pending_id` +- `data/logs/debug/YYYY-MM-DD.log` — DEBUG sink, gzipped, 30-day retention +- `data/logs/info/YYYY-MM-DD.log` — INFO sink, gzipped, 30-day retention +- `JobStore` is in-memory only today; Stage 3 will move it to SQLite. **Key conventions:** - All modules include the file header `# coding=utf-8` plus a `@Time` / `@Author` docstring. Match this when adding new files. -- `common/logging.py` configures a single process-wide `loguru` logger and removes the default handler. Import `from common.logging import logger` rather than creating new loggers. +- `common/logging.py` configures a single process-wide `loguru` logger (stderr + two rotated file sinks). Import `from common.logging import logger` rather than creating new loggers. - `spark_executor/__init__.py` re-exports `app` so callers can `from spark_executor import app` — keep this re-export when adding to the package. - The `FastApiMCP` instance 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-mcp` passes tool args as a JSON body, and dict-typed query params arrive as strings and 422. Every new MCP tool needs a request model in `tools/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 `KeyError` for "unknown id" and `ValueError` for 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_job` snapshots the connection's `master` / `deploy_mode` / `spark_conf` into the `PendingSubmission`; `confirm_submit_job` is the only place `spark-submit` is 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`) in `monkeypatch.setattr` because the tool modules captured the originals at import time. See `tests/integration/test_mcp_routes.py` for the pattern. +- **86 tests pass** as of the last Stage 1 cleanup; run `uv run pytest` after 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_file` for LLM-written PySpark code (`job_writer` + one new tool). +- **Stage 3: deferred** — async submission, SQLite `JobRegistry`, status poller, offset-based `LogCache`, multi-tenant `owner`. Plan lives at `docs/superpowers/plans/2026-06-24-spark-executor-mcp.md`.