From 0323fd36d45abad13f03601f03a9f101367f350f Mon Sep 17 00:00:00 2001 From: "J. Nick Koston" Date: Wed, 7 Oct 2026 10:00:07 -1000 Subject: [PATCH] [core] Clone git PlatformIO packages in parallel in the prefetch (#18837) --- esphome/platformio/prefetch.py | 79 +++++++++---- tests/unit_tests/test_platformio_prefetch.py | 114 ++++++++++++++++++- 2 files changed, 166 insertions(+), 27 deletions(-) diff --git a/esphome/platformio/prefetch.py b/esphome/platformio/prefetch.py index fbb31ae452..991164eb99 100644 --- a/esphome/platformio/prefetch.py +++ b/esphome/platformio/prefetch.py @@ -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 diff --git a/tests/unit_tests/test_platformio_prefetch.py b/tests/unit_tests/test_platformio_prefetch.py index ef573767f7..3f01fb22e6 100644 --- a/tests/unit_tests/test_platformio_prefetch.py +++ b/tests/unit_tests/test_platformio_prefetch.py @@ -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+")