[core] Clone git PlatformIO packages in parallel in the prefetch (#18837)

This commit is contained in:
J. Nick Koston
2026-10-08 09:00:07 +13:00
committed by GitHub
parent 1ba323e9a3
commit 0323fd36d4
2 changed files with 166 additions and 27 deletions
+56 -23
View File
@@ -62,6 +62,9 @@ def _preserved_sys_path() -> Iterator[None]:
# Concurrent registry resolutions / HEAD probes (each is network-bound)
_RESOLVE_WORKERS = 8
# Concurrent VCS clones; network-bound, so wider than the extraction pool
_CLONE_WORKERS = 8
# A hung child must not block the build; downloads resume on the next run
_PREFETCH_TIMEOUT = 20 * 60
@@ -372,42 +375,64 @@ def _registry_jobs(
return jobs, failed, installable
def _is_vcs_spec_uri(url: str) -> bool:
"""Whether pio's ``install_from_uri`` would clone this URI rather than
copy or download it (PackageSpec normalizes git URLs to ``git+``)."""
return not url.startswith(("file://", "symlink://", "http://", "https://"))
def _is_clone_entry(entry: tuple) -> bool:
"""Whether this pre-install entry's spec is cloned, not extracted."""
url = entry[1].uri
return bool(url) and _is_vcs_spec_uri(url)
def _spec_name(spec: Any, url: str) -> str:
"""The spec's name; the URL basename fallback is defensive only
(PackageSpec derives a name from the URI itself)."""
return spec.name or url.split("#", 1)[0].rsplit("/", 1)[-1]
def _uri_jobs(
manager: Any, specs: list[Any], seen: set[str]
manager: Any, specs: list[Any], seen: set[str], trusted_names: bool = False
) -> tuple[list[tuple[str, int, Any]], int, list[tuple[str, Any]]]:
"""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) and the ``(name, spec)`` pairs whose archives will be
installable.
error) and the ``(name, spec)`` pairs to pre-install, including VCS
specs, which the pre-install clones itself. ``trusted_names`` marks
platform packages, whose platform.json keys match their manifests.
"""
from esphome.net_retry import fetch_with_retry, http_request
candidates: list[tuple[str, str, Path, Any]] = []
candidates: list[tuple[str, str, Path, Any, bool]] = []
installable: list[tuple[str, Any]] = []
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 not url:
continue
is_vcs = _is_vcs_spec_uri(url)
if not is_vcs and not url.startswith(("http://", "https://")):
continue # file/symlink specs are copied in place by pio run
if manager.get_package(spec):
continue
name = spec.name or url.rsplit("/", 1)[-1]
name = _spec_name(spec, url)
# Only a name that is also the install dir may pre-install
safe_name = trusted_names or spec.has_custom_name()
if is_vcs:
if safe_name:
installable.append((name, spec))
continue
# PlatformIO downloads URL specs with no checksum
dl_path = Path(manager.compute_download_path(url, ""))
if dl_path.is_file():
if spec.has_custom_name():
# Only a custom name (Foo=https://...) is the destination
# dir; a URI-derived name's destination comes from the
# archive manifest, so its dedupe key could collide with
# another name and race one directory. pio run installs it.
if safe_name:
installable.append((name, spec)) # fetched by an earlier run
continue
if str(dl_path) in seen:
continue # another spec already claimed this .part
seen.add(str(dl_path))
candidates.append((spec.name, url, dl_path, spec))
candidates.append((name, url, dl_path, spec, safe_name))
errors: list[str] = []
@@ -434,16 +459,17 @@ def _uri_jobs(
if not candidates:
return [], 0, installable
with ThreadPoolExecutor(max_workers=min(_RESOLVE_WORKERS, len(candidates))) as ex:
sizes = list(ex.map(_head_size, [url for _, url, _, _ in candidates]))
sizes = list(ex.map(_head_size, [url for _, url, _, _, _ in candidates]))
jobs: list[tuple[str, int, Any]] = []
failed = 0
for (name, url, dl_path, spec), size in zip(candidates, sizes, strict=True):
for (name, url, dl_path, spec, safe_name), 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)))
if spec.has_custom_name():
# See above: derived-name specs stay with pio run's installer
if safe_name:
installable.append((name, spec))
else:
# Missing or unusable Content-Length; visible under -v
@@ -703,7 +729,11 @@ def _preinstall(
would hang, not fail). Waves skip dependencies; the installed
manifests feed the next wave. Any failure falls back to pio run.
"""
workers = extract_workers(len(entries))
# Clones wait on the network: sort them first, dependency waves too
entries = sorted(entries, key=lambda entry: not _is_clone_entry(entry))
clones = sum(map(_is_clone_entry, entries))
# Network-bound clones get a wider pool than CPU-bound extraction
workers = max(extract_workers(len(entries)), min(clones, _CLONE_WORKERS))
# One manager per worker (_install mutates instance state); built
# serially because construction rewires the shared manager logger
managers: SimpleQueue = SimpleQueue()
@@ -735,8 +765,9 @@ def _preinstall(
raise
_LOGGER.info(
"Installing %d PlatformIO package(s) with %d extraction worker(s): %s",
"Installing %d PlatformIO package(s)%s with %d worker(s): %s",
len(entries),
f" ({clones} clone(s))" if clones else "",
workers,
", ".join(name for name, *_ in entries),
)
@@ -871,8 +902,10 @@ def _prefetch(build_dir: Path, env: str) -> None:
unresolved = 0
for mgr, batch, is_platform in ((p.pm, specs, True), (lm, lib_specs, False)):
entries: list[tuple[str, Any]] = []
for build_jobs in (_registry_jobs, _uri_jobs):
batch_jobs, failed, installable = build_jobs(mgr, batch, seen)
for batch_jobs, failed, installable in (
_registry_jobs(mgr, batch, seen),
_uri_jobs(mgr, batch, seen, trusted_names=is_platform),
):
jobs += batch_jobs
unresolved += failed
entries += installable
+110 -4
View File
@@ -661,7 +661,8 @@ def test_registry_jobs_one_bad_spec_keeps_the_rest(tmp_path: Path) -> None:
def test_uri_jobs_head_sizes_the_bar(tmp_path: Path) -> None:
"""HEAD sizes direct-URL specs; git and unreachable URLs are skipped."""
"""HEAD sizes direct-URL specs; VCS specs skip the download but are
still installable (the pre-install clones them in parallel)."""
m = _fake_manager(tmp_path)
resp = MagicMock()
resp.headers = {"content-length": "2222"}
@@ -670,15 +671,24 @@ def test_uri_jobs_head_sizes_the_bar(tmp_path: Path) -> None:
m,
[
_FakeSpec(uri="https://x/big.zip", name="big", custom_name=True),
_FakeSpec(uri="git+https://x/repo.git", name="repo"),
_FakeSpec(uri="https://x/repo.git#v1", name="barevcs"),
_FakeSpec(uri="git+https://x/repo.git", name="repo", custom_name=True),
_FakeSpec(name="registry"),
],
set(),
)
assert failed == 0
assert [(n, s) for n, s, _ in jobs] == [("big", 2222)]
assert [n for n, _ in installable] == ["big"]
assert [n for n, _ in installable] == ["repo", "big"]
# A derived-name platform archive sized in the same run is trusted too
with patch("esphome.net_retry.http_request", return_value=resp):
jobs, _, installable = pf._uri_jobs(
m,
[_FakeSpec(uri="https://x/fresh.zip", name="fresh")],
set(),
trusted_names=True,
)
assert [n for n, _, _ in jobs] == ["fresh"]
assert [n for n, _ in installable] == ["fresh"]
# a successful HEAD with no Content-Length is a clean skip
resp.headers = {}
with patch("esphome.net_retry.http_request", return_value=resp):
@@ -687,6 +697,53 @@ def test_uri_jobs_head_sizes_the_bar(tmp_path: Path) -> None:
) == ([], 0, [])
def test_uri_jobs_vcs_specs_installable_without_probe(tmp_path: Path) -> None:
"""VCS specs never probe the network here (there is no archive); an
uninstalled custom-named one is handed to the pre-install, while
derived names, installed specs, and file/symlink specs are skipped."""
m = _fake_manager(tmp_path)
with patch("esphome.net_retry.http_request") as mock_head:
jobs, failed, installable = pf._uri_jobs(
m,
[
_FakeSpec(
uri="git+https://x/tool.git#1.0", name="tool", custom_name=True
),
_FakeSpec(uri="hg+https://x/old", name="mercurial", custom_name=True),
# A derived name is not the destination dir; pio run clones it
_FakeSpec(uri="git+https://x/derived#v2", name="derived"),
_FakeSpec(uri="file:///local/dir", name="local"),
_FakeSpec(uri="symlink:///local/dir", name="link"),
],
set(),
)
mock_head.assert_not_called()
assert (jobs, failed) == ([], 0)
assert [n for n, _ in installable] == ["tool", "mercurial"]
# trusted_names admits platform packages, which have no custom name
with patch("esphome.net_retry.http_request"):
_, _, installable = pf._uri_jobs(
m,
[_FakeSpec(uri="git+https://x/derived#v2", name="derived")],
set(),
trusted_names=True,
)
assert [n for n, _ in installable] == ["derived"]
m.get_package.return_value = object() # already installed: warm and silent
with patch("esphome.net_retry.http_request"):
assert pf._uri_jobs(
m,
[
_FakeSpec(
uri="git+https://x/tool.git#1.0", name="tool", custom_name=True
)
],
set(),
) == ([], 0, [])
def test_uri_jobs_head_failure_counts_as_unresolved(
tmp_path: Path, caplog: pytest.LogCaptureFixture
) -> None:
@@ -750,6 +807,11 @@ def test_uri_jobs_skips_installed_cached_and_seen(tmp_path: Path) -> None:
jobs, failed, installable = pf._uri_jobs(m, spec, set())
assert (jobs, failed) == ([], 0)
assert [n for n, _ in installable] == ["a"]
# a cached platform archive with a derived name is trusted the same way
derived = [_FakeSpec(uri="https://x/a.zip", name="a")]
assert pf._uri_jobs(m, derived, set())[2] == []
_, _, installable = pf._uri_jobs(m, derived, set(), trusted_names=True)
assert [n for n, _ in installable] == ["a"]
dl.unlink()
# a registry job already claimed this download path
assert pf._uri_jobs(m, spec, {str(dl)}) == ([], 0, [])
@@ -875,6 +937,44 @@ def test_dependency_entries_filter_seen_names(tmp_path: Path) -> None:
assert pf._dependency_entries(m, [("top@1", _FakeSpec(name="top"))], {"dep"}) == []
def test_preinstall_widens_the_pool_for_clones(
tmp_path: Path, monkeypatch: pytest.MonkeyPatch
) -> None:
"""Network-bound clones all run at once even when the CPU-sized
extraction pool would serialize them; the barrier deadlocks otherwise."""
m = _fake_manager(tmp_path)
monkeypatch.setattr(pf, "extract_workers", lambda jobs=None: 1)
barrier = threading.Barrier(3)
m._install.side_effect = lambda spec, **kw: barrier.wait(timeout=5)
pf._preinstall(
m,
[
(f"g{i}", _FakeSpec(uri=f"git+https://x/g{i}.git", name=f"g{i}"))
for i in range(3)
],
)
assert barrier.broken is False
def test_preinstall_orders_clones_before_extractions(
tmp_path: Path, monkeypatch: pytest.MonkeyPatch
) -> None:
"""A clone must not queue behind CPU-bound archive extractions."""
m = _fake_manager(tmp_path)
monkeypatch.setattr(pf, "extract_workers", lambda jobs=None: 1)
calls: list[str] = []
m._install.side_effect = lambda spec, **kw: calls.append(spec.name)
pf._preinstall(
m,
[
("a@1", _FakeSpec(name="a")),
("b@1", _FakeSpec(name="b")),
("g@1", _FakeSpec(uri="git+https://x/g.git", name="g")),
],
)
assert calls[0] == "g"
def test_preinstall_cleanup_cannot_displace_the_inflight_error(
tmp_path: Path, caplog: pytest.LogCaptureFixture, monkeypatch: pytest.MonkeyPatch
) -> None:
@@ -1890,3 +1990,9 @@ def test_platformio_private_api_contract() -> None:
derived = PackageSpec("https://x/y/archive/master.zip")
assert derived.name and not derived.has_custom_name()
assert PackageSpec("Foo=https://x/y/archive/master.zip").has_custom_name()
# _is_vcs_spec_uri relies on bare .git URLs normalizing to git+, on
# both parse paths (raw string, and requirements= for platform tools)
assert PackageSpec("https://github.com/x/y.git#v1").uri.startswith("git+")
assert PackageSpec(
owner="o", name="tool-x", requirements="https://github.com/x/y.git"
).uri.startswith("git+")