Add a new optional field to Connection for the Spark History Server
base URL. The field is stored as part of the connection but is not
yet consumed by any tool — for now it's a labeled place to record
where SHS lives on the cluster, and a hook for future SHS-specific
tools. To actually fetch SHS endpoints today, use fetch_url with the
SHS host added to url_allowlist.
The field is Optional[str], default None. Same nullability as
yarn_rm_url. The SHS host and the YARN RM host are usually
different, so this is independent of yarn_rm_url.
Schema:
- models.py: Connection.history_server_url (str | None, default None)
- requests.py:
- SaveConnectionRequest.history_server_url (str | None, default None)
- UpdateConnectionRequest.history_server_url (str | None, default None)
- tools/connections.py: save_connection signature gains
history_server_url with the _UNSET sentinel pattern (same as
yarn_rm_url, url_allowlist, etc.) so the upsert path correctly
distinguishes "not provided" from "explicitly None".
- tools/connections.py: update_connection docstring lists the new
field in the mutable fields set.
Tests (4 new in tests/unit/test_connection_tools.py):
- test_save_connection_with_history_server_url
- test_save_connection_history_server_url_defaults_to_none
- test_update_connection_changes_history_server_url
- test_update_connection_keeps_history_server_url_when_omitted
Tests: 398 passed (was 394, +4 net).
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
166 lines
5.4 KiB
Python
166 lines
5.4 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'."
|
|
),
|
|
)
|
|
history_server_url: str | None = Field(
|
|
default=None,
|
|
description=(
|
|
"Optional Spark History Server (SHS) base URL, e.g. "
|
|
"'http://history.prod.internal:18080'. Stored as part of the "
|
|
"connection for reference and to be consumed by future SHS-"
|
|
"specific tools. Currently **not consumed by any tool** — to "
|
|
"fetch SHS endpoints today, use fetch_url with the SHS host "
|
|
"added to url_allowlist. The SHS host and the YARN RM host "
|
|
"are usually different, so this is independent of yarn_rm_url."
|
|
),
|
|
)
|
|
|
|
@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
|