tao.chenandClaude 44c9f415b7 tools: 修 spark_submit 校验的两处小问题
c9e54a3 引入了 server-minted 校验, 但有两处需要调整:

1. 错误信息用 fmt.Errorf(... %w).Error() 构造, %w 包装语义被丢
   掉 (errResult 接收 string, Error() 之后 %w 已经不可达). 改用
   fmt.Sprintf 拼接, %q 引用路径, 错误信息保持不变但代码不再误
   导.

2. 校验放在 scriptPath 解析之后、queue 解析之前. 错误信息会按字
   段出现顺序报 (cluster_id, master, deploy_mode, script_path,
   queue, ...), 但当前顺序是 script_path 校验先报, 然后才报 queue
   缺失. 把校验挪到所有 RequireString 之后、parseStringMap 之前,
   LLM 看错误时字段顺序跟 schema 顺序一致.

nil-safe 校验逻辑不变 (d.UploadStore == nil 时跳过). 现有测试
全部通过:
  - TestSparkSubmit_RejectsNonMintedPath 仍通过 (校验位置不影响
    行为)
  - TestSparkSubmit_StructuredCommand 仍通过 (走 mint 路径)

Co-Authored-By: Claude <noreply@anthropic.com>
2026-07-13 19:55:16 +08:00

spark-mcp-go

MCP (Model Context Protocol) server for Apache Spark on YARN. Lets an LLM agent discover Spark/YARN endpoints, submit jobs, fetch logs, and analyze them via 11 typed Tools. Streamable HTTP transport, SQLite-backed configuration, and log/slog structured logging with per-Tool call files.

What it does

  • Exposes 11 MCP Tools over Streamable HTTP at /mcp
  • Admin API at /admin/* for cluster and audit configuration
  • HTTP Basic and YARN SimpleAuth for Spark/YARN endpoints
  • SSRF protection with cluster URL allowlist plus DNS rebinding guard
  • Per-Tool independent audit log (file mode 0600)
  • spark-submit via local exec.Command (no shell, no injection)

Quick start (3 steps)

  1. Copy and edit the environment file:

    cp .env.example .env
    # Edit ADMIN_TOKENS and AGENT_TOKEN
    
  2. Build and start the server:

    go build ./...
    ./spark-mcp-go
    # Or use the helper:
    # ./scripts/dev.sh
    
  3. Initialize an MCP session, then call list_clusters to discover endpoints.

11 Tools

Tool Type Purpose
list_clusters discovery Return all active clusters (LLM entry point)
spark_submit exec Local spark-submit process, extracts app_id
list_applications RM GET /ws/v1/cluster/apps[?state=&user=]
get_application_status RM GET /ws/v1/cluster/apps/{id}
get_application_logs RM Fallback chain: amContainerLogs -> aggregated-logs -> logs
kill_application RM PUT state KILLED
fetch_spark_metrics SHS Executor metrics + summary mode triggers analyzer
fetch_cluster_env RM Aggregate /cluster/info and /cluster/metrics
analyze_spark_log analyzer RM logs + heuristic rules + LLM-ready prompt
fetch_url HTTP Generic primitive, uses cluster allowlist and auth
upload_file FS Write to ./data/uploads/, reference from spark_submit

Environment variables

Variable Default Required Description
LISTEN_ADDR :8080 no HTTP listen address
DATA_DIR ./data no SQLite, uploads, and log root
ADMIN_TOKENS yes Comma-separated tokens for /admin/*
AGENT_TOKEN yes Token for /mcp requests
HTTP_CLIENT_TIMEOUT 30s no Timeout for fetch_url/RM Tool HTTP calls
MAX_RESPONSE_BYTES 1048576 (1 MiB) no HTTP response truncation limit
SPARK_SUBMIT_TIMEOUT 60s no spark-submit process timeout
LOG_DIR ./data/logs no Log root directory
LOG_LEVEL info no debug, info, warn, or error
LOG_FORMAT text no text for terminal, json for log files
ANALYZER_DATA_SKEW_RATIO 3.0 no Data-skew rule threshold (max/min ratio)
ANALYZER_GC_PRESSURE_RATIO 0.1 no GC-pressure rule threshold (GC/CPU ratio)
ANALYZER_BOTTLENECK_SHUFFLE_GB 50 no Bottleneck rule threshold (shuffle GB)

Admin API

All endpoints require Authorization: Bearer <admin_token>.

Method Path Description
GET /admin/clusters List all clusters
POST /admin/clusters Create a cluster
GET /admin/clusters/:id Get one cluster
PUT /admin/clusters/:id Update a cluster
DELETE /admin/clusters/:id Delete a cluster
GET /admin/audit?limit=100 Query audit log

Create a cluster:

curl -sS -X POST \\
  -H "Authorization: Bearer ${ADMIN_TOKEN}" \\
  -H "Content-Type: application/json" \\
  -d '{
    "id": "prod",
    "name": "prod",
    "rm_url": "http://rm.example.com:8088",
    "shs_url": "http://shs.example.com:18080",
    "spark_submit_execute_bin": "/opt/spark/bin/spark-submit",
    "is_active": true,
    "auth_type": "simple",
    "auth_username": "yarn",
    "rate_limit_per_min": 10,
    "url_allowlist": ["rm.example.com:8088", "shs.example.com:18080"]
  }' \\
  "http://127.0.0.1:${LISTEN_ADDR:-:8080}/admin/clusters"

Note: auth_password is not accepted via JSON (json:"-"). Set it through a dedicated password endpoint or seed the database directly.

Security model

Three lines of defense:

  1. Token auth - ADMIN_TOKENS protects /admin/*; AGENT_TOKEN protects /mcp.
  2. SSRF guard - Every outbound URL must match the cluster's url_allowlist, with private-IP CIDR blacklist (incl. 169.254.0.0/16 AWS/GCP metadata).
  3. Path allowlist - upload_file writes only to ./data/uploads/, rejects absolute paths, .., and non-[a-zA-Z0-9._-] filenames.

Additional guarantees:

  • spark_submit runs the cluster's binary via exec.Command(name, args...) — slice form, never sh -c. No shell metacharacter interpretation.
  • auth_password is tagged json:"-"; never serialized in responses, never accepted from JSON request bodies.
  • Cross-host redirects (RM → NM 307) preserve the Authorization header through DoWithRedirect only when the destination host is in url_allowlist.
  • Per-Tool log files are created with mode 0600 and live under ./data/logs/tools/.

Development

go build ./...
go test ./...
go vet ./...

More docs

S
Description
No description provided
Readme
268 KiB
2026-07-15 16:37:20 +08:00
Languages
Go 86.5%
HTML 12.6%
Shell 0.9%