Files
mcp-server/spark_executor/models.py
T
ClaudeandClaude Fable 5 f43e5b6403 docs(fetch_url): fix 3 misleading description bits
Three small but high-leverage text corrections. No code or test
changes — these only affect what the LLM sees in tools/list and
Pydantic schemas.

1) FetchUrlRequest.connection_name description
   Old: "The Connection's yarn_rm_url defines the allowed host domain."
   New: explicitly says url_allowlist is the host gate, NOT yarn_rm_url.
        yarn_rm_url is only used by get_external_* tools. Without this
        fix the LLM would try to control fetch scope via yarn_rm_url
        (a no-op) instead of url_allowlist.

2) Connection.url_allowlist description
   Added: "Set or change via save_connection (pass url_allowlist on
   create) or update_connection (PATCH the field on an existing
   connection)." Tells the LLM which tools populate the field,
   instead of leaving it to guess.

3) /fetch_url route description
   Old: "30s timeout, redirects followed."
   New: "30s timeout, redirects followed, response body capped at
        1 MB (the response includes a truncated boolean when this
        kicks in)." LLM previously had no way to discover the 1MB
        cap or the truncated indicator; it would just see
        short responses and assume that was the full body.

Tests: 394 passed, no test changes (descriptions are not asserted).

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-07-09 16:18:29 +08:00

154 lines
4.8 KiB
Python

# coding=utf-8
"""
@Time :2026/6/24
@Author :tao.chen
"""
from datetime import datetime
from pydantic import BaseModel, Field, field_validator
class Job(BaseModel):
job_id: str
application_id: str
script_path: str
queue: str
submit_time: datetime
connection: str
yarn_rm_url: str | None = None
class JobStatus(BaseModel):
application_id: str
state: str
raw: str = Field(default="")
class JobResult(BaseModel):
application_id: str
state: str
final_status: str | None = None
diagnostics: str | None = None
tracking_url: str | None = None
started_time: int | None = None
finished_time: int | None = None
class SubmitResult(BaseModel):
job_id: str
application_id: str
tracking_url: str | None = None
class FetchUrlResult(BaseModel):
url: str
status_code: int
content_type: str
body: str
class ApplicationSummary(BaseModel):
"""A YARN application summary from /ws/v1/cluster/apps.
Field names are mapped from the YARN JSON keys to clearer
snake_case names by the tool function. Unused YARN fields
(memorySeconds, vcoreSeconds, preemptedResource*, etc.) are
not exposed — the LLM doesn't need them.
"""
application_id: str
name: str
user: str
queue: str
state: str
final_status: str | None = None
application_type: str | None = None
application_tags: str = ""
started_time: int = 0
finished_time: int = 0
tracking_url: str | None = None
progress: float | None = None
class Connection(BaseModel):
name: str
# Defaults to "yarn" because that's the literal string spark-submit wants
# for --master when targeting YARN. Override for Standalone (spark://...),
# Kubernetes (k8s://...), or local mode.
master: str = "yarn"
deploy_mode: str = "cluster"
yarn_rm_url: str | None = None
spark_conf: dict[str, str] = Field(default_factory=dict)
# None means "fall back to Settings.ssl_verify_default". Explicit True/False
# overrides the global default for this connection.
ssl_verify: bool | None = None
ssl_ca_bundle: str | None = None
# Authentication for YARN REST calls.
auth_type: str = "none" # "none" | "simple" | "basic" | "kerberos"
auth_user: str | None = None
auth_password: str | None = None
# Display/audit only for kerberos; actual SPNEGO uses the system cache.
auth_principal: str | None = None
auth_keytab: str | None = None
url_allowlist: list[str] = Field(
default_factory=list,
description=(
"List of fnmatch glob patterns for hosts the fetch_url tool may access. "
"The list is mandatory-opt-in: an empty list (the default) denies all "
"hosts, so you must populate it before fetch_url can access any URL. "
"Set or change via save_connection (pass url_allowlist on create) or "
"update_connection (PATCH the field on an existing connection). "
"Useful for clusters whose hostnames do NOT share a common suffix — "
"e.g. single-label hosts like 'ccam1'-'ccam99' (configure ['ccam*']) "
"or HDFS namenode on a different subdomain ('*.hadoop.internal'). "
"Patterns are matched against the URL host only (no port, no path). "
"fnmatch rules apply: '*' does NOT match '.', so 'ccam*' matches "
"'ccam50' but not 'ccam50.evil.com'."
),
)
@field_validator("master")
@classmethod
def _check_master(cls, v: str) -> str:
"""Catch common typos like 'yarn-cluster' or 'http://...'. """
if v == "yarn":
return v
if v.startswith(("spark://", "k8s://", "mesos://", "local")):
return v
raise ValueError(
f"master must be 'yarn', 'spark://...', 'k8s://...', 'mesos://...', "
f"or 'local[/N]'; got {v!r}"
)
@field_validator("auth_type")
@classmethod
def _check_auth_type(cls, v: str) -> str:
if v not in {"none", "simple", "basic", "kerberos"}:
raise ValueError(
f"auth_type must be one of none/simple/basic/kerberos; got {v!r}"
)
return v
class PendingSubmission(BaseModel):
pending_id: str
app_name: str | None = None
connection: str
master: str
deploy_mode: str
yarn_rm_url: str | None = None
script_path: str
queue: str
executor_memory: str
executor_cores: int
num_executors: int
spark_conf: dict[str, str] = Field(default_factory=dict)
extra_args: dict[str, str] = Field(default_factory=dict)
created_at: datetime
status: str = "PENDING" # PENDING | SUBMITTED | CANCELLED | FAILED
error: str | None = None
job_id: str | None = None
application_id: str | None = None
tracking_url: str | None = None