[core] Prefetch PlatformIO packages in parallel (#18769)

This commit is contained in:
J. Nick Koston
2026-08-27 14:42:30 +12:00
committed by GitHub
parent e53870085b
commit 2104096f02
9 changed files with 1713 additions and 49 deletions
+9 -19
View File
@@ -2,7 +2,6 @@
from collections.abc import Callable
from ctypes.util import find_library
from functools import partial
import json
import logging
import os
@@ -24,16 +23,17 @@ from esphome.framework_helpers import (
create_venv,
download_and_extract,
download_from_mirrors,
download_with_resume,
failure_reason,
get_python_env_executable_path,
get_system_python_path,
resume_fetch_job,
rmdir,
run_batch_downloads,
run_command,
run_command_ok,
str_to_lst_of_str,
tool_version_runs,
warn_prefetch_failures,
)
from esphome.helpers import write_file_if_changed
@@ -686,18 +686,6 @@ def _patch_tools_json_demote_unused_tools(framework_path: Path) -> None:
)
def _download_tool(
dist_path: Path, entry: dict, tracker: Callable[[int], None]
) -> None:
download_with_resume(
entry["url"],
dist_path / entry["dest"],
sha256=entry["sha256"],
size=entry["size"],
progress=tracker,
)
def _prefetch_idf_tool_archives(
framework_path: Path,
targets_str: str,
@@ -775,15 +763,17 @@ def _prefetch_idf_tool_archives(
(
entry["name"],
entry["size"],
partial(_download_tool, dist_path, entry),
resume_fetch_job(
entry["url"],
dist_path / entry["dest"],
sha256=entry["sha256"],
size=entry["size"],
),
)
for entry in entries
],
)
for name, e in failures:
# failure_reason: a message-less exception must not log blank
_LOGGER.warning("Could not prefetch %s: %s", name, failure_reason(e))
_LOGGER.debug("Prefetch failure detail", exc_info=e)
warn_prefetch_failures(failures)
if len(failures) == len(entries):
# A systematic fault, not one flaky mirror: the resume
# workaround (#17703) is off for this whole install
+41 -3
View File
@@ -701,7 +701,7 @@ def _write_download_meta(
_LOGGER.debug("Could not update download metadata %s: %s", meta, e)
def _content_length(resp: "requests.Response") -> int:
def content_length(resp: "requests.Response") -> int:
"""Return the response's Content-Length, or 0 when absent or malformed.
0 means "unknown", which downstream disables the progress bar and the
@@ -744,7 +744,7 @@ def _stream_response_to_file(
"""
f.seek(offset)
f.truncate(offset)
total_size = size or offset + _content_length(resp)
total_size = size or offset + content_length(resp)
downloaded = offset
own_bar: ProgressBar | None = None
if progress is None:
@@ -909,6 +909,19 @@ def _part_path(dest: Path) -> Path:
return dest.with_name(dest.name + ".part")
def discard_partial_download(dest: Path) -> None:
"""Remove ``dest`` and the resume sidecars of an abandoned download."""
part = _part_path(dest)
for stale in (dest, part, part.with_name(part.name + ".meta")):
try:
stale.unlink()
except FileNotFoundError:
continue
except OSError as err:
# The caller's cache is never pruned; leave a trace
_LOGGER.debug("Could not remove %s: %s", stale, err)
def _cancellable_sleep(
delay: float, progress: Callable[[int], None] | None, done: int
) -> None:
@@ -922,6 +935,31 @@ def _cancellable_sleep(
time.sleep(min(0.5, remaining))
def resume_fetch_job(
url: str, dest: PathType, **kwargs
) -> Callable[[Callable[[int], None]], None]:
"""A ``run_batch_downloads`` job callable wrapping ``download_with_resume``.
Forwards the runner's positional tracker as the ``progress`` keyword.
"""
def fetch(tracker: Callable[[int], None]) -> None:
download_with_resume(url, dest, progress=tracker, **kwargs)
return fetch
def warn_prefetch_failures(
failures: list[tuple[str, BaseException]],
message: str = "Could not prefetch %s: %s",
) -> None:
"""Warn per failed batch-prefetch job; the caller's installer retries them."""
for name, err in failures:
# failure_reason: a message-less exception must not log blank
_LOGGER.warning(message, name, failure_reason(err))
_LOGGER.debug("Prefetch failure detail", exc_info=err)
def download_with_resume(
url: str,
dest: PathType,
@@ -1022,7 +1060,7 @@ def download_with_resume(
streamed = True
if offset == 0:
validator = _response_validator(resp)
expected_total = _content_length(resp)
expected_total = content_length(resp)
# Recorded so a later run can prove an If-Range
# resume of this part file safe.
_write_download_meta(meta, url, validator, expected_total)
+5 -8
View File
@@ -35,6 +35,7 @@ from esphome.framework_helpers import (
failure_reason,
rmdir,
run_batch_downloads,
warn_prefetch_failures,
)
_LOGGER = logging.getLogger(__name__)
@@ -977,14 +978,10 @@ def _prefetch_wave(
for c in components
],
)
for name, err in failures:
# The sequential call below retries and raises the real error
_LOGGER.warning(
"Prefetch of %s failed (retrying sequentially): %s",
name,
failure_reason(err),
)
_LOGGER.debug("Prefetch failure detail", exc_info=err)
# The sequential call below retries and raises the real error
warn_prefetch_failures(
failures, "Prefetch of %s failed (retrying sequentially): %s"
)
except Exception as err: # noqa: BLE001 # pylint: disable=broad-exception-caught
# Same policy as the ESP-IDF twin: the prefetch must never become a
# new way for the build to fail
+600
View File
@@ -0,0 +1,600 @@
"""Parallel prefetch of the packages a PlatformIO run would install.
Downloads the archives concurrently into PlatformIO's own download cache
(identical ``compute_download_path`` keys) so the serial installer finds
them already cached. Runs in a subprocess like all PlatformIO execution:
loading a platform executes its code (pioarduino's penv setup rewrites
``sys.path``). A sentinel in the build dir lets warm builds skip the
spawn. Best-effort: any failure logs and PlatformIO downloads as before.
Across processes sharing a core dir every download destination is
serialized by a file lock; checksum-less URL downloads additionally
stage under a stable name and promote with an atomic rename.
"""
from __future__ import annotations
from concurrent.futures import ThreadPoolExecutor
import hashlib
import json
import logging
import os
from pathlib import Path
import subprocess
import sys
import threading
import time
from typing import Any
from esphome.framework_helpers import (
content_length,
discard_partial_download,
failure_reason,
resume_fetch_job,
run_batch_downloads,
warn_prefetch_failures,
)
from esphome.helpers import get_bool_env
_LOGGER = logging.getLogger(__name__)
# Concurrent registry resolutions / HEAD probes (each is network-bound)
_RESOLVE_WORKERS = 8
# A hung child must not block the build; downloads resume on the next run
_PREFETCH_TIMEOUT = 20 * 60
# Waiting on another process's URL download; past this, leave it to pio
_DOWNLOAD_LOCK_TIMEOUT = 60
# Child exit for a handled, already-warned failure; 1 would collide with
# the interpreter's own import-failure exit
_EXIT_HANDLED = 3
# Short lock-acquire slices so a waiting worker still observes Ctrl-C
_URI_LOCK_POLL = 1
# Resolution errored (vs a clean skip); suppresses the warm sentinel
_RESOLVE_FAILED = object()
def _sweep_stale_sidecars(download_dir: Path, expire_seconds: int) -> None:
"""Prune resume sidecars pio's usage.db pruner cannot see.
A version bump strands an aborted archive's sidecars forever. Lock
files stay: a held lock can carry an ancient mtime (O_TRUNC keeps
it), and unlinking one reopens the single-writer hole it guards.
"""
cutoff = time.time() - expire_seconds
try:
for f in download_dir.iterdir():
if f.suffix not in (".part", ".meta", ".prefetch"):
continue
try:
if f.stat().st_mtime < cutoff:
f.unlink()
except OSError as err:
_LOGGER.debug("Could not remove %s: %s", f, err)
except OSError:
_LOGGER.debug("Could not sweep %s", download_dir, exc_info=True)
# Child records a no-work run; the parent skips the next spawn while valid
_SENTINEL_NAME = ".esphome_prefetch.json"
_SENTINEL_SCHEMA = 1
def _ini_sha256(build_dir: Path) -> str:
return hashlib.sha256((build_dir / "platformio.ini").read_bytes()).hexdigest()
def _sentinel_state(build_dir: Path) -> dict[str, Any]:
"""The environment fingerprint a sentinel must match to stay valid."""
# Same fingerprint as the heal stamp: the sentinel's dirs die with its wipe
from esphome.platformio.toolchain import current_python_minor
return {
"schema": _SENTINEL_SCHEMA,
"ini_sha256": _ini_sha256(build_dir),
"python": current_python_minor(),
"core_dir_env": os.environ.get("PLATFORMIO_CORE_DIR", ""),
}
def _prefetch_is_warm(build_dir: Path) -> bool:
"""Whether the last prefetch found nothing to do and nothing changed since."""
try:
data = json.loads((build_dir / _SENTINEL_NAME).read_text(encoding="utf-8"))
dirs = data.pop("dirs")
return (
data == _sentinel_state(build_dir)
and bool(dirs)
and all(Path(d).is_dir() for d in dirs)
)
except FileNotFoundError:
return False
except (OSError, ValueError, KeyError, AttributeError, TypeError):
_LOGGER.debug("Ignoring invalid prefetch sentinel", exc_info=True)
return False
def prefetch_platformio_packages() -> None:
"""Warm PlatformIO's download cache for the current project, in parallel."""
from esphome.core import CORE
from esphome.platformio.toolchain import (
default_libdeps_dir,
heal_platformio_python_env,
)
# Heal first: its Python-version wipe would discard freshly warmed
# caches and the sentinel's dirs (the later heal call is a no-op)
heal_platformio_python_env()
build_dir = Path(CORE.build_path)
if _prefetch_is_warm(build_dir):
return
# The child is esphome itself: PYTHONPATH stays so it imports this
# tree's esphome (tests/integration pins the source tree through it)
env = dict(os.environ)
# Must match run_platformio_cli's default or warm builds re-resolve
# every library
env.setdefault("PLATFORMIO_LIBDEPS_DIR", default_libdeps_dir())
# -v/-vv must reach the child's debug logging or the swallowed
# failure detail is undiagnosable in the field
env["ESPHOME_PREFETCH_LOG_LEVEL"] = str(logging.getLogger().getEffectiveLevel())
if CORE.dashboard:
# The child's progress bar and log escaping key off CORE.dashboard
env["ESPHOME_PREFETCH_DASHBOARD"] = "1"
cmd = [
sys.executable,
"-m",
"esphome.platformio.prefetch",
str(build_dir),
CORE.name,
]
try:
proc = subprocess.run(cmd, env=env, check=False, timeout=_PREFETCH_TIMEOUT)
except subprocess.TimeoutExpired:
_LOGGER.warning("PlatformIO package prefetch timed out; continuing without it")
return
except Exception as err: # noqa: BLE001 # pylint: disable=broad-exception-caught
# The prefetch must never become a new way for the build to fail
_LOGGER.warning("PlatformIO package prefetch skipped: %s", failure_reason(err))
_LOGGER.debug("Prefetch failure detail", exc_info=True)
return
if proc.returncode == _EXIT_HANDLED:
# The child already warned with the reason; a second line is noise
_LOGGER.debug("Prefetch child reported a handled failure")
elif proc.returncode != 0:
# Exit 1 stays here: the interpreter exits 1 for import/module
# failures before main() ever runs, a wiring break worth a warning
_LOGGER.warning(
"PlatformIO package prefetch skipped (exit %d)", proc.returncode
)
def _project_platform_and_config(ini: Path, env: str) -> tuple[str | None, Any]:
"""The env's platform spec and the ProjectConfig for the given ini."""
from platformio import app
from platformio.project.config import ProjectConfig
# PlatformBase.config reads the default ProjectConfig; it must see
# this ini's env options
app.set_session_var("custom_project_conf", str(ini))
config = ProjectConfig.get_instance(str(ini))
return config.get(f"env:{env}", "platform", None), config
def _registry_jobs(
manager, specs, seen: set[str]
) -> tuple[list[tuple[str, int, Any]], int]:
"""Resolve registry specs to ``(name, size, fetch)`` batch jobs.
Mirrors PlatformIO's install path: best version, systype file, first
mirror, and the same sha1(url + checksum) download-cache key. Also
returns how many resolutions errored (a clean skip is not an error).
"""
from platformio.registry.mirror import RegistryFileMirrorIterator
local = threading.local()
errors: list[str] = []
def _resolve(spec) -> tuple[str, int, str, Path, str] | object | None:
# One manager (and registry HTTP session) per worker thread;
# installed-state was already checked on the shared manager
if (mgr := getattr(local, "mgr", None)) is None:
mgr = local.mgr = manager.__class__()
try:
packages = mgr.search_registry_packages(spec)
if not packages:
_LOGGER.debug("%s is unknown to the registry", spec)
return None # let PlatformIO report it
package, version = mgr.find_best_registry_version(packages, spec)
if not package or not version:
_LOGGER.debug("%s has no matching registry version", spec)
return None
pkgfile = mgr.pick_compatible_pkg_file(version["files"])
if not pkgfile:
_LOGGER.debug("%s has no file for this systype", spec)
return None
url, checksum = next(RegistryFileMirrorIterator(pkgfile["download_url"]))
checksum = checksum or pkgfile["checksum"]["sha256"]
dl_path = Path(mgr.compute_download_path(url, checksum))
if dl_path.is_file():
return None # cached from an earlier run
size = pkgfile.get("size")
if not size:
_LOGGER.debug("%s has no size; PlatformIO fetches it", spec)
return None # no size, no bar share
return f"{package['name']}@{version['name']}", size, url, dl_path, checksum
except Exception as err: # noqa: BLE001 # pylint: disable=broad-exception-caught
# One flaky spec must not discard the rest of the batch
_LOGGER.debug("Could not resolve %s", spec, exc_info=True)
errors.append(failure_reason(err))
return _RESOLVE_FAILED
# Serial disk lookups on the shared manager: a fully warm build
# resolves nothing, and duplicate specs resolve once
unique: dict[tuple[str | None, str, str], Any] = {}
for s in specs:
if not s.uri and not manager.get_package(s):
unique.setdefault((s.owner, s.name, str(s.requirements)), s)
pending = list(unique.values())
if not pending:
return [], 0
# Serial resolutions (registry GET + mirror HEAD each) dominate
with ThreadPoolExecutor(max_workers=min(_RESOLVE_WORKERS, len(pending))) as ex:
results = list(ex.map(_resolve, pending))
jobs: list[tuple[str, int, Any]] = []
for res in results:
if res is None or res is _RESOLVE_FAILED:
continue
name, size, url, dl_path, checksum = res
if str(dl_path) in seen:
continue # duplicate spec; two workers must not share a .part
seen.add(str(dl_path))
jobs.append(
(name, size, _registry_fetch_job(manager, url, dl_path, checksum, size))
)
if failed := len(errors):
# Visible once per build, naming a cause so an API break does not
# read as an outage; per-spec detail stays at debug
_LOGGER.warning(
"Could not resolve %d of %d PlatformIO package(s) (%s); "
"PlatformIO will download them serially",
failed,
len(pending),
errors[0],
)
return jobs, failed
def _uri_jobs(manager, specs, seen: set[str]) -> tuple[list[tuple[str, int, Any]], int]:
"""Jobs for direct-URL specs; a HEAD sizes each for the combined bar.
Also returns how many HEAD probes errored (an absent length is not an
error).
"""
from esphome.net_retry import fetch_with_retry, http_request
candidates: list[tuple[str, str, Path]] = []
for spec in specs:
url = spec.uri
if not url or not url.startswith(("http://", "https://")):
continue # git+/file specs are cloned/copied, not downloaded
if url.split("#", 1)[0].endswith(".git"):
continue # bare-URL VCS spec; PlatformIO clones it
if manager.get_package(spec):
continue
# PlatformIO downloads URL specs with no checksum
dl_path = Path(manager.compute_download_path(url, ""))
if dl_path.is_file() or str(dl_path) in seen:
continue # cached, or another spec already claimed this .part
seen.add(str(dl_path))
candidates.append((spec.name, url, dl_path))
errors: list[str] = []
def _head_size(url: str) -> int:
try:
resp = fetch_with_retry(url, lambda: http_request("HEAD", url, timeout=30))
except Exception as err: # noqa: BLE001 # pylint: disable=broad-exception-caught
_LOGGER.debug("HEAD %s failed", url, exc_info=True)
errors.append(failure_reason(err))
return -1
if not resp.ok:
# An error page's Content-Length is not a download size
_LOGGER.debug("HEAD %s returned %s", url, resp.status_code)
if resp.status_code in (401, 403, 408, 429) or resp.status_code >= 500:
# 401/403 included: registries rate-limit with them
errors.append(f"HTTP {resp.status_code}")
return -1 # transient; must not be cached as warm
# Permanent (405/501 HEAD-unsupported, 401/403/404): a clean
# skip so the warm sentinel is not disabled forever; pio run
# surfaces a genuinely broken URL when it downloads
return 0
return content_length(resp)
if not candidates:
return [], 0
with ThreadPoolExecutor(max_workers=min(_RESOLVE_WORKERS, len(candidates))) as ex:
sizes = list(ex.map(_head_size, [url for _, url, _ in candidates]))
jobs: list[tuple[str, int, Any]] = []
failed = 0
for (name, url, dl_path), size in zip(candidates, sizes, strict=True):
if size < 0:
failed += 1
elif size:
jobs.append((name, size, _uri_fetch_job(manager, url, dl_path, size)))
else:
# Missing or unusable Content-Length; visible under -v
_LOGGER.debug("%s reports no usable length; PlatformIO fetches it", url)
if failed:
_LOGGER.warning(
"Could not size %d of %d PlatformIO package URL(s) (%s); "
"PlatformIO will download them serially",
failed,
len(candidates),
errors[0],
)
return jobs, failed
def _serialized_fetch_job(
dl_path: Path, lock_path: str, body: Any, unlocked_ok: bool = True
) -> Any:
"""Wrap ``body`` so the shared destination is single-writer.
Interleaved writers truncate each other's ``.part`` bytes (see
registry.py). The bounded poll observes Ctrl-C via the tracker; a
blown deadline is a clean skip (the holder's copy is what the build
needs). On a lock-less filesystem a sha256-verified body runs
unlocked with one warning; a checksum-less one
(``unlocked_ok=False``) is a counted failure instead.
"""
def run(tracker: Any) -> None:
from filelock import FileLock, Timeout
# fallback_to_soft would leave a stale marker on lock-less
# filesystems that blocks every later build (see git.py)
lock = FileLock(lock_path, fallback_to_soft=False)
deadline = time.monotonic() + _DOWNLOAD_LOCK_TIMEOUT
while True:
try:
lock.acquire(timeout=_URI_LOCK_POLL)
break
except Timeout:
tracker(0) # raises when the batch is cancelled
if time.monotonic() >= deadline:
# Another process is fetching this same file; its copy
# is what the build needs (a large framework archive
# can hold the lock far longer than this deadline)
_LOGGER.debug("Leaving %s to its current downloader", dl_path.name)
return
except OSError as err:
if not unlocked_ok:
# A body with no checksum to catch interleaved corruption
raise
lock = None
_LOGGER.warning(
"Could not lock %s (%s); downloading unlocked",
dl_path.name,
err,
)
break
try:
if dl_path.is_file():
return # another process finished it while we waited
body(tracker)
finally:
if lock is not None:
lock.release()
return run
# usage.db is a whole-file rewrite behind pio's self-unlinking LockFile;
# concurrent writers could reset every recorded entry
_REGISTER_LOCK = threading.Lock()
def _register_download(manager: Any, dl_path: Path) -> None:
"""Hand the archive to pio's usage.db pruner; an unregistered one is
never expired (disk garbage, never a bad build)."""
try:
with _REGISTER_LOCK:
manager.set_download_utime(str(dl_path))
except Exception as err: # noqa: BLE001 # pylint: disable=broad-exception-caught
_LOGGER.debug("Could not register %s with pio's cache: %s", dl_path, err)
def _registry_fetch_job(
manager: Any, url: str, dl_path: Path, checksum: str, size: int
) -> Any:
"""A locked fetch straight to the cache path; sha256 verifies it."""
# .esphome.lock: pio's own LockFile(dl_path) owns <dl_path>.lock and
# deletes it on release, which would unlink a held filelock
fetch = _serialized_fetch_job(
dl_path,
f"{dl_path}.esphome.lock",
resume_fetch_job(url, dl_path, sha256=checksum, size=size),
)
def run(tracker: Any) -> None:
fetch(tracker)
if dl_path.is_file():
# The deadline skip can end with no archive landed
_register_download(manager, dl_path)
return run
def _uri_fetch_job(manager: Any, url: str, dl_path: Path, size: int) -> Any:
"""Fetch to a locked staging path, then rename into the cache.
The stable staging name keeps resume working across interrupted
runs; the rename makes the promotion atomic.
"""
tmp = dl_path.with_name(f"{dl_path.name}.prefetch")
# attempts=2: the size is only a HEAD probe's word, and a HEAD/GET
# disagreement would otherwise re-download the archive five times
fetch = resume_fetch_job(url, tmp, size=size, attempts=2)
def promote(tracker: Any) -> None:
fetch(tracker)
if (actual := tmp.stat().st_size) != size:
# A wrong-length checksum-less body must never be published
discard_partial_download(tmp)
raise ValueError(f"expected {size} bytes, fetched {actual}")
tmp.replace(dl_path)
def run(tracker: Any) -> None:
_serialized_fetch_job(dl_path, f"{tmp}.lock", promote, unlocked_ok=False)(
tracker
)
if dl_path.is_file():
# Won or lost, the race is over; staging files left behind
# are dead weight PlatformIO's cache never prunes
discard_partial_download(tmp)
_register_download(manager, dl_path)
return run
def _prefetch(build_dir: Path, env: str) -> None:
from platformio.dependencies import get_core_dependencies
from platformio.package.manager.library import LibraryPackageManager
from platformio.package.manager.platform import PlatformPackageManager
from platformio.package.meta import PackageSpec
from platformio.platform.factory import PlatformFactory
platform_spec, config = _project_platform_and_config(
build_dir / "platformio.ini", env
)
if not platform_spec:
# An env mismatch must not disable the feature with no trace
_LOGGER.debug(
"No platform for env %s in %s; nothing to prefetch", env, build_dir
)
return
# The platform (manifest plus build scripts) installs first and
# resolves the rest. Its setup may rewrite sys.path (pioarduino's penv
# setup does); restore it so later imports here still resolve.
saved_sys_path = list(sys.path)
pm = PlatformPackageManager()
_sweep_stale_sidecars(Path(pm.get_download_dir()), pm.DOWNLOAD_CACHE_EXPIRE)
pkg = pm.install(platform_spec, skip_dependencies=True)
p = PlatformFactory.new(pkg)
p.configure_project_packages(env, ["run"])
sys.path[:] = saved_sys_path
specs = [
p.get_package_spec(name)
for name, opts in p.packages.items()
if not opts.get("optional")
]
# PIO's build engine installs outside the platform package list;
# skipped when the platform lists it itself
if not any(s.name == "tool-scons" for s in specs):
specs.append(
PackageSpec(
owner="platformio",
name="tool-scons",
requirements=get_core_dependencies()["tool-scons"],
)
)
lib_deps = config.get(f"env:{env}", "lib_deps", [])
# pio run's storage dir for this env: installed libraries skip by
# disk lookup
libdeps_dir = Path(config.get("platformio", "libdeps_dir")) / env
lm = LibraryPackageManager(str(libdeps_dir))
# A bare name is usually a framework built-in (WiFi, SPI); with no
# lib builders here to tell built-in from registry, skip it. The only
# cost is that an owner-less user library is not prefetched
lib_specs = [
spec
for dep in lib_deps
if dep and not dep.startswith("$")
if (spec := PackageSpec(dep)).external or spec.owner
]
seen: set[str] = set()
jobs: list[tuple[str, int, Any]] = []
unresolved = 0
for mgr, batch in ((p.pm, specs), (lm, lib_specs)):
for build_jobs in (_registry_jobs, _uri_jobs):
batch_jobs, failed = build_jobs(mgr, batch, seen)
jobs += batch_jobs
unresolved += failed
sentinel = build_dir / _SENTINEL_NAME
if not jobs:
if not unresolved:
# Record the no-work run so the parent skips the next spawn.
# A failed resolution is not "no work": a registry outage must
# not be cached as warm.
dirs = [config.get("platformio", "packages_dir")]
if lib_specs:
dirs.append(str(libdeps_dir))
sentinel.write_text(
json.dumps({**_sentinel_state(build_dir), "dirs": dirs}),
encoding="utf-8",
)
return
sentinel.unlink(missing_ok=True)
_LOGGER.info(
"Prefetching %d PlatformIO package(s): %s",
len(jobs),
", ".join(name for name, _, _ in jobs),
)
# PlatformIO retries failed packages itself, without resume
warn_prefetch_failures(run_batch_downloads("Downloading PlatformIO packages", jobs))
def main(argv: list[str]) -> int:
"""Subprocess entry point: ``prefetch <build_dir> <env_name>``."""
from esphome.core import CORE
from esphome.log import setup_log
raw_level = os.environ.get("ESPHOME_PREFETCH_LOG_LEVEL")
try:
level = int(raw_level) if raw_level is not None else logging.INFO
except ValueError:
level = logging.INFO
# Mirror the parent's log setup: warnings keep their level prefix and
# color, and the download bar still draws under the dashboard
CORE.dashboard = get_bool_env("ESPHOME_PREFETCH_DASHBOARD")
setup_log(level)
# pio's managers attach their own handler and still propagate; without
# this every manager line also prints through the root handler. Their
# construction re-pins the logger to INFO, so a logger-level filter
# (which survives pio's handler reset) enforces a quiet level instead.
for cls_name in (
"ToolPackageManager",
"LibraryPackageManager",
"PlatformPackageManager",
):
manager_logger = logging.getLogger(cls_name.replace("Package", " "))
manager_logger.propagate = False
manager_logger.addFilter(lambda record: record.levelno >= level)
if len(argv) != 2:
# A wiring bug, not a network failure; make it distinguishable
_LOGGER.warning("prefetch usage: <build_dir> <env_name>")
return 2
build_dir, env = argv
try:
_prefetch(Path(build_dir), env)
except KeyboardInterrupt:
# Shared process group: exit quietly, no traceback on the terminal
_LOGGER.debug("Prefetch interrupted", exc_info=True)
return 130
except Exception as err: # noqa: BLE001 # pylint: disable=broad-exception-caught
# The parent treats any exit as warn-and-continue, never a failure
_LOGGER.warning("PlatformIO package prefetch skipped: %s", failure_reason(err))
_LOGGER.debug("Prefetch failure detail", exc_info=True)
return _EXIT_HANDLED
return 0
if __name__ == "__main__": # pragma: no cover
sys.exit(main(sys.argv[1:]))
+12 -5
View File
@@ -96,7 +96,7 @@ def _clean_platformio_python_env(config: "ProjectConfig", core_dir: Path) -> Non
rmtree(penv)
def _current_python_minor() -> str:
def current_python_minor() -> str:
"""Return the running interpreter's ``major.minor`` (e.g. ``3.13``)."""
return f"{sys.version_info.major}.{sys.version_info.minor}"
@@ -161,7 +161,7 @@ def heal_platformio_python_env() -> None:
def _check_platformio_python_stamp(config: "ProjectConfig") -> None:
"""Compare the stamp to the running interpreter; wipe and restamp on mismatch."""
current = _current_python_minor()
current = current_python_minor()
stamp_dir = _pio_stamp_dir(config)
# Host the stamp/lock even before PlatformIO's first run creates the dir.
stamp_dir.mkdir(parents=True, exist_ok=True)
@@ -289,6 +289,12 @@ def copy_ccache_script() -> None:
)
def default_libdeps_dir() -> str:
"""The PLATFORMIO_LIBDEPS_DIR value a pio run defaults to; the package
prefetch must resolve installed libraries against the same dir."""
return str(CORE.relative_piolibdeps_path().absolute())
def run_platformio_cli(*args, **kwargs) -> str | int:
# Re-provision the PlatformIO cache if the interpreter's major.minor changed
# since it was last built; a stale platform otherwise rejects the new Python
@@ -296,9 +302,7 @@ def run_platformio_cli(*args, **kwargs) -> str | int:
heal_platformio_python_env()
os.environ["PLATFORMIO_FORCE_COLOR"] = "true"
os.environ["PLATFORMIO_BUILD_DIR"] = str(CORE.relative_pioenvs_path().absolute())
os.environ.setdefault(
"PLATFORMIO_LIBDEPS_DIR", str(CORE.relative_piolibdeps_path().absolute())
)
os.environ.setdefault("PLATFORMIO_LIBDEPS_DIR", default_libdeps_dir())
# Suppress Python syntax warnings from third-party scripts during compilation
os.environ.setdefault("PYTHONWARNINGS", "ignore::SyntaxWarning")
# Increase uv retry count to handle transient network errors (default is 3)
@@ -346,6 +350,9 @@ def run_platformio_cli_run(config, verbose, *args, **kwargs) -> str | int:
def run_compile(config, verbose):
from esphome.platformio.prefetch import prefetch_platformio_packages
prefetch_platformio_packages()
args = []
if CONF_COMPILE_PROCESS_LIMIT in config[CONF_ESPHOME]:
args += [f"-j{config[CONF_ESPHOME][CONF_COMPILE_PROCESS_LIMIT]}"]