fix(yarn_client): fall back to directory listing when /stdout is HTML

Some YARN deployments (Knox-proxied, wrapped vendor builds) return an
HTML directory listing page for /node/containerlogs/.../stdout instead of
the raw log bytes. The previous fix assumed a clean text/plain response
and failed every log fetch on those clusters.

Two new private helpers in yarn_client.py:
  - _looks_like_html(text): cheap prolog sniff (<!doctype html, <html,
    <?xml) on the first ~200 chars.
  - _parse_log_listing_html(html): extract plain log-file names via a
    href="..." regex; drop ../, absolute paths, sort toggles (?C=N;O=D),
    and absolute URLs.

_fetch_logs_via_am_container now does:
  1. GET /apps/{appid} (unchanged).
  2. Read app.amContainerLogs.
  3. Fast path: GET {amContainerLogs}/stdout. If 2xx and NOT HTML, return.
  4. Fallback: GET {amContainerLogs} (bare URL), parse the HTML listing,
     fetch each log file individually, and concatenate with === filename ===
     headers in document order.
  5. Raise _logs_unavailable_error if both paths yield nothing.

Tests: 3 new (test_logs_falls_back_to_directory_listing_when_stdout_returns_html,
test_logs_html_filter_strips_sort_links_and_parent_dir,
test_logs_html_filter_handles_prelaunch_files,
test_logs_raises_when_directory_listing_empty) — full suite 247 passed.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
Claude
2026-06-30 10:56:42 +08:00
co-authored by Claude Fable 5
parent fd0036ca80
commit 32ee167b59
2 changed files with 206 additions and 10 deletions
+72 -10
View File
@@ -14,9 +14,12 @@ Endpoints used (YARN 2.6+):
When aggregated logs are unavailable (HTTP 404/501, e.g. log aggregation When aggregated logs are unavailable (HTTP 404/501, e.g. log aggregation
disabled), fall back to the amContainerLogs field reported by the app disabled), fall back to the amContainerLogs field reported by the app
endpoint and fetch the AM (driver) container's /stdout directly from the endpoint and fetch the AM (driver) container's /stdout from the
NodeManager. This returns driver logs only; executor container logs still NodeManager. Some clusters (Knox-proxied, wrapped YARN builds) wrap
require yarn.log-aggregation-enable=true (the primary /aggregated-logs path). /stdout in an HTML directory listing; in that case we fetch the bare
{amContainerLogs} URL, parse the listing, and pull each log file. This
returns driver logs only; executor container logs still require
yarn.log-aggregation-enable=true (the primary /aggregated-logs path).
The ResourceManager URL is passed in per call (snapshotted on the Job at The ResourceManager URL is passed in per call (snapshotted on the Job at
confirm_submit_job time) and falls back to the YARN_RESOURCE_MANAGER_URL env confirm_submit_job time) and falls back to the YARN_RESOURCE_MANAGER_URL env
@@ -26,6 +29,8 @@ field was designed for, but no longer requires the `yarn` CLI to interpret it.
import json import json
from dataclasses import dataclass from dataclasses import dataclass
import re
import httpx import httpx
import httpx_kerberos import httpx_kerberos
@@ -163,11 +168,50 @@ def _logs_unavailable_error(application_id: str) -> YarnError:
) )
def _looks_like_html(text: str) -> bool:
"""Cheap heuristic: True if the body starts with an HTML/XHTML prolog.
Intentionally lightweight (no full parser) — used only to decide whether
a NodeManager log response is raw text or an HTML directory listing.
"""
head = text.lstrip()[:200].lower()
return head.startswith(("<!doctype html", "<html", "<?xml"))
def _parse_log_listing_html(html: str) -> list[str]:
"""Extract plain log-file names from a NodeManager directory-listing page.
Keeps simple filenames (stdout, stderr, syslog, prelaunch.err,
launch_container.sh, ...); drops parent-dir links, absolute URLs, and
sort-toggle query strings. Returns hrefs in document order, deduplicated.
"""
hrefs = re.findall(r'href="([^"]*)"', html, re.IGNORECASE)
seen: set[str] = set()
out: list[str] = []
for href in hrefs:
if href in ("../", "/") or href.startswith("/"):
continue
if "://" in href or "?" in href:
continue
if href in seen:
continue
seen.add(href)
out.append(href)
return out
def _fetch_logs_via_am_container(application_id: str, config: YarnClientConfig) -> str: def _fetch_logs_via_am_container(application_id: str, config: YarnClientConfig) -> str:
""" """
Fetch ApplicationMaster (driver) container stdout as a fallback when the Fetch ApplicationMaster (driver) container logs as a fallback when the
aggregated-logs endpoint is unavailable. aggregated-logs endpoint is unavailable.
Fast path: GET {amContainerLogs}/stdout and return its text directly when
the NodeManager serves raw log bytes (modern YARN).
Fallback path: when /stdout returns an HTML directory listing (Knox /
wrapped YARN), GET the bare {amContainerLogs} URL, parse the listing for
log-file names, and fetch each one individually.
Returns AM (driver) container logs only. Executor container logs require Returns AM (driver) container logs only. Executor container logs require
yarn.log-aggregation-enable=true, which is covered by the primary yarn.log-aggregation-enable=true, which is covered by the primary
/aggregated-logs path. /aggregated-logs path.
@@ -183,12 +227,30 @@ def _fetch_logs_via_am_container(application_id: str, config: YarnClientConfig)
if not am_container_logs: if not am_container_logs:
raise _logs_unavailable_error(application_id) raise _logs_unavailable_error(application_id)
log_url = f"{am_container_logs.rstrip('/')}/stdout" am_container_logs = am_container_logs.rstrip("/")
resp = _request("GET", log_url, timeout=60.0, verify=config.verify_for_httpx(), auth=config.auth_for_httpx()) verify = config.verify_for_httpx()
if resp.status_code >= 400: auth = config.auth_for_httpx()
logger.error(f"YARN GET {log_url} -> {resp.status_code}: {resp.text[:500]}")
raise _logs_unavailable_error(application_id) # Fast path: modern YARN serves raw log bytes at /stdout.
return resp.text stdout_url = f"{am_container_logs}/stdout"
resp = _request("GET", stdout_url, timeout=60.0, verify=verify, auth=auth)
if resp.status_code < 400 and not _looks_like_html(resp.text):
return resp.text
# Fallback path: Knox / wrapped YARN serves an HTML directory listing for
# every URL on the NodeManager log endpoint. Parse it and fetch each file.
resp = _request("GET", am_container_logs, timeout=60.0, verify=verify, auth=auth)
if resp.status_code < 400:
parts: list[str] = []
for filename in _parse_log_listing_html(resp.text):
file_url = f"{am_container_logs}/{filename}"
file_resp = _request("GET", file_url, timeout=60.0, verify=verify, auth=auth)
if file_resp.status_code < 400 and not _looks_like_html(file_resp.text):
parts.append(f"=== {filename} ===\n{file_resp.text}")
if parts:
return "\n\n".join(parts)
raise _logs_unavailable_error(application_id)
def get_application_logs(application_id: str, config: YarnClientConfig) -> str: def get_application_logs(application_id: str, config: YarnClientConfig) -> str:
+134
View File
@@ -12,6 +12,7 @@ from spark_executor.core.yarn_client import (
YarnConfigError, YarnConfigError,
YarnError, YarnError,
YarnClientConfig, YarnClientConfig,
_parse_log_listing_html,
get_application_logs, get_application_logs,
get_application_status, get_application_status,
kill_application, kill_application,
@@ -201,6 +202,139 @@ def test_logs_5xx_on_aggregated_endpoint_raises_immediately():
assert m.call_count == 1 assert m.call_count == 1
def test_logs_falls_back_to_directory_listing_when_stdout_returns_html():
"""Knox / wrapped YARN: GET {amContainerLogs}/stdout returns an HTML
directory listing instead of raw log bytes. The fallback should fetch the
bare URL, parse the listing, and pull each log file individually."""
listing = (
"<html><body>"
'<a href="../">Parent Directory</a>'
'<a href="?C=N;O=D">Name</a>'
'<a href="stdout">stdout</a>'
'<a href="stderr">stderr</a>'
'<a href="syslog">syslog</a>'
"</body></html>"
)
html_stdout = (
"<!doctype html><html><body>wrapper around /stdout endpoint</body></html>"
)
responses = [
_resp(404), # aggregated-logs missing
_resp(
200,
json_data={
"app": {
"amContainerLogs": "http://nm1:8042/node/containerlogs/container_1/hdfs",
}
},
),
_resp(200, text=html_stdout), # /stdout -> HTML wrapper
_resp(200, text=listing), # bare URL -> directory listing
_resp(200, text="driver stdout line\n"), # fetch stdout
_resp(200, text="driver stderr line\n"), # fetch stderr
_resp(200, text="syslog line\n"), # fetch syslog
]
with patch(
"spark_executor.core.yarn_client.httpx.request", side_effect=responses
) as m:
out = get_application_logs("application_1", YarnClientConfig(yarn_rm_url=RM))
# Header + body for each file in document order.
assert "=== stdout ===" in out
assert "=== stderr ===" in out
assert "=== syslog ===" in out
assert "driver stdout line" in out
assert "driver stderr line" in out
assert "syslog line" in out
# stdout must come before stderr (document order preserved).
assert out.index("=== stdout ===") < out.index("=== stderr ===") < out.index("=== syslog ===")
# 7 HTTP calls: aggregated-logs, app, /stdout (HTML), bare URL, stdout file, stderr file, syslog file
assert m.call_count == 7
assert m.call_args_list[2].args == (
"GET",
"http://nm1:8042/node/containerlogs/container_1/hdfs/stdout",
)
assert m.call_args_list[3].args == (
"GET",
"http://nm1:8042/node/containerlogs/container_1/hdfs",
)
assert m.call_args_list[4].args == (
"GET",
"http://nm1:8042/node/containerlogs/container_1/hdfs/stdout",
)
assert m.call_args_list[5].args == (
"GET",
"http://nm1:8042/node/containerlogs/container_1/hdfs/stderr",
)
assert m.call_args_list[6].args == (
"GET",
"http://nm1:8042/node/containerlogs/container_1/hdfs/syslog",
)
def test_logs_html_filter_strips_sort_links_and_parent_dir():
"""_parse_log_listing_html must drop ../, sort-toggle query strings, and
absolute URLs, while keeping plain log-file names."""
html = (
'<a href="../">Parent</a>'
'<a href="?C=N;O=D">Name</a>'
'<a href="?C=M;O=A">Last modified</a>'
'<a href="?C=S;O=A">Size</a>'
'<a href="?C=D;O=A">Description</a>'
'<a href="http://example.com/somewhere">absolute</a>'
'<a href="/absolute/path">rooted</a>'
'<a href="stdout">stdout</a>'
'<a href="stderr">stderr</a>'
'<a href="syslog">syslog</a>'
'<a href="stdout">stdout-dup</a>' # duplicate
)
assert _parse_log_listing_html(html) == ["stdout", "stderr", "syslog"]
def test_logs_html_filter_handles_prelaunch_files():
"""Some YARN versions list prelaunch.out / prelaunch.err / launch_container.sh
in the directory listing alongside stdout/stderr. Keep them — we want to
surface every log file the container produced."""
html = (
'<a href="prelaunch.out">prelaunch.out</a>'
'<a href="prelaunch.err">prelaunch.err</a>'
'<a href="launch_container.sh">launch_container.sh</a>'
'<a href="directory-info">directory-info</a>'
)
assert _parse_log_listing_html(html) == [
"prelaunch.out",
"prelaunch.err",
"launch_container.sh",
"directory-info",
]
def test_logs_raises_when_directory_listing_empty():
"""If /stdout returns HTML AND the directory listing has no parseable log
file names, raise _logs_unavailable_error."""
html_stdout = "<!doctype html><html>wrapper</html>"
empty_listing = '<html><body><a href="../">Parent</a></body></html>'
responses = [
_resp(404), # aggregated-logs missing
_resp(
200,
json_data={
"app": {
"amContainerLogs": "http://nm1:8042/node/containerlogs/container_1/hdfs",
}
},
),
_resp(200, text=html_stdout), # /stdout -> HTML
_resp(200, text=empty_listing), # bare URL -> only parent link
]
with patch(
"spark_executor.core.yarn_client.httpx.request", side_effect=responses
):
with pytest.raises(YarnError, match="log-aggregation-enable"):
get_application_logs("application_1", YarnClientConfig(yarn_rm_url=RM))
# --- kill_application --- # --- kill_application ---
def test_kill_sends_put_with_killed_state(): def test_kill_sends_put_with_killed_state():