diff --git a/.github/workflows/ci-artifact-transport.yml b/.github/workflows/ci-artifact-transport.yml index 8b937cc1aef5..1ed374adfc9b 100644 --- a/.github/workflows/ci-artifact-transport.yml +++ b/.github/workflows/ci-artifact-transport.yml @@ -9,6 +9,7 @@ on: - workers/ci-artifacts/** - scripts/ci/restore-r2-artifact.py - scripts/ci/node_product_cache.py + - scripts/ci/peer_product_source.py - scripts/ci/restore-app-host-test-product.sh - tests/test_node_product_cache.py - scripts/ci/verify-r2-canary.py @@ -23,6 +24,7 @@ on: - workers/ci-artifacts/** - scripts/ci/restore-r2-artifact.py - scripts/ci/node_product_cache.py + - scripts/ci/peer_product_source.py - scripts/ci/restore-app-host-test-product.sh - tests/test_node_product_cache.py - scripts/ci/verify-r2-canary.py diff --git a/.github/workflows/ci-macos.yml b/.github/workflows/ci-macos.yml index 2af92a21834a..8cde4edfe947 100644 --- a/.github/workflows/ci-macos.yml +++ b/.github/workflows/ci-macos.yml @@ -58,10 +58,13 @@ jobs: product_contract: ${{ steps.product-key.outputs.key }} source_revision: ${{ github.sha }} producer_run_id: ${{ github.run_id }} + producer_run_attempt: ${{ github.run_attempt }} env: CMUX_NODE_PRODUCT_CACHE_ROOT: ${{ vars.CMUX_NODE_PRODUCT_CACHE_ROOT }} CMUX_NODE_PRODUCT_CACHE_MAX_BYTES: ${{ vars.CMUX_NODE_PRODUCT_CACHE_MAX_BYTES }} CMUX_NODE_PRODUCT_CACHE_WAIT_SECONDS: ${{ vars.CMUX_NODE_PRODUCT_CACHE_WAIT_SECONDS }} + CMUX_ARTIFACT_PEER_URLS: ${{ vars.CMUX_ARTIFACT_PEER_URLS }} + CMUX_ARTIFACT_PEER_TOKEN_FILE: ${{ vars.CMUX_ARTIFACT_PEER_TOKEN_FILE }} CMUX_CI_XCODE_APP: ${{ vars.CMUX_CI_XCODE_APP_MACOS_15 }} CMUX_CI_REQUIRED_MACOS_SDK_MAJOR: "26" CMUX_SKIP_ZIG_BUILD: "1" @@ -712,6 +715,7 @@ jobs: CMUX_PRODUCT_CONTRACT: ${{ steps.product-key.outputs.key }} CMUX_PRODUCT_SOURCE_REVISION: ${{ github.sha }} CMUX_PRODUCT_PRODUCER_RUN_ID: ${{ github.run_id }} + CMUX_PRODUCT_PRODUCER_RUN_ATTEMPT: ${{ github.run_attempt }} run: python3 scripts/ci/node_product_cache.py seed "$RUNNER_TEMP/app-host-products.tar.gz" app-host-unit-tests: @@ -744,6 +748,8 @@ jobs: CMUX_NODE_PRODUCT_CACHE_ROOT: ${{ vars.CMUX_NODE_PRODUCT_CACHE_ROOT }} CMUX_NODE_PRODUCT_CACHE_MAX_BYTES: ${{ vars.CMUX_NODE_PRODUCT_CACHE_MAX_BYTES }} CMUX_NODE_PRODUCT_CACHE_WAIT_SECONDS: ${{ vars.CMUX_NODE_PRODUCT_CACHE_WAIT_SECONDS }} + CMUX_ARTIFACT_PEER_URLS: ${{ vars.CMUX_ARTIFACT_PEER_URLS }} + CMUX_ARTIFACT_PEER_TOKEN_FILE: ${{ vars.CMUX_ARTIFACT_PEER_TOKEN_FILE }} # This independent job-level marker makes every app-host wrapper fail # closed if a setup step or environment handoff loses either redirect. CMUX_CI_APP_HOST_ISOLATION_REQUIRED: "1" @@ -890,12 +896,29 @@ jobs: CMUX_PRODUCT_CONTRACT: ${{ needs.macos-compile-admission.outputs.product_contract }} CMUX_PRODUCT_SOURCE_REVISION: ${{ needs.macos-compile-admission.outputs.source_revision }} CMUX_PRODUCT_PRODUCER_RUN_ID: ${{ needs.macos-compile-admission.outputs.producer_run_id }} - CMUX_NODE_PRODUCT_CACHE_FALLBACK_SOURCE: ${{ vars.CI_ARTIFACT_R2_URL != '' && 'r2' || 'github' }} + CMUX_PRODUCT_PRODUCER_RUN_ATTEMPT: ${{ needs.macos-compile-admission.outputs.producer_run_attempt }} + CMUX_NODE_PRODUCT_CACHE_FALLBACK_SOURCE: ${{ vars.CMUX_ARTIFACT_PEER_URLS != '' && 'peer' || (vars.CI_ARTIFACT_R2_URL != '' && 'r2' || 'github') }} run: python3 scripts/ci/node_product_cache.py acquire "$RUNNER_TEMP/app-host-products" + - name: Try trusted fleet peer artifact source + id: peer-products + if: steps.node-products.outputs.hit != 'true' + continue-on-error: true + env: + ARTIFACT_ID: ${{ needs.macos-compile-admission.outputs.artifact_id }} + ARTIFACT_PROVIDER_DIGEST: ${{ needs.macos-compile-admission.outputs.artifact_digest }} + EXPECTED_SHA256: ${{ needs.macos-compile-admission.outputs.sha256 }} + CMUX_PRODUCT_CONTRACT: ${{ needs.macos-compile-admission.outputs.product_contract }} + CMUX_PRODUCT_SOURCE_REVISION: ${{ needs.macos-compile-admission.outputs.source_revision }} + CMUX_PRODUCT_PRODUCER_RUN_ID: ${{ needs.macos-compile-admission.outputs.producer_run_id }} + CMUX_PRODUCT_PRODUCER_RUN_ATTEMPT: ${{ needs.macos-compile-admission.outputs.producer_run_attempt }} + run: python3 scripts/ci/peer_product_source.py fetch "$RUNNER_TEMP/app-host-products" + + # R2 is an optional remote broker. With no configured endpoint, a peer + # miss falls straight through to the canonical GitHub artifact. - name: Try shared R2 artifact transport id: r2-products - if: steps.node-products.outputs.hit != 'true' + if: steps.node-products.outputs.hit != 'true' && steps.peer-products.outputs.hit != 'true' && vars.CI_ARTIFACT_R2_URL != '' continue-on-error: true env: GH_TOKEN: ${{ github.token }} @@ -905,7 +928,7 @@ jobs: run: python3 scripts/ci/restore-r2-artifact.py - name: Download compiled app-host test product - if: steps.node-products.outputs.hit != 'true' && steps.r2-products.outputs.hit != 'true' + if: steps.node-products.outputs.hit != 'true' && steps.peer-products.outputs.hit != 'true' && steps.r2-products.outputs.hit != 'true' uses: ./.github/actions/download-test-product with: artifact-id: ${{ needs.macos-compile-admission.outputs.artifact_id }} @@ -914,9 +937,20 @@ jobs: - name: Restore compiled app-host test product id: restore-products env: + ARTIFACT_ID: ${{ needs.macos-compile-admission.outputs.artifact_id }} + ARTIFACT_PROVIDER_DIGEST: ${{ needs.macos-compile-admission.outputs.artifact_digest }} EXPECTED_SHA256: ${{ needs.macos-compile-admission.outputs.sha256 }} + CMUX_PRODUCT_CONTRACT: ${{ needs.macos-compile-admission.outputs.product_contract }} + CMUX_PRODUCT_SOURCE_REVISION: ${{ needs.macos-compile-admission.outputs.source_revision }} + CMUX_PRODUCT_PRODUCER_RUN_ID: ${{ needs.macos-compile-admission.outputs.producer_run_id }} + CMUX_PRODUCT_PRODUCER_RUN_ATTEMPT: ${{ needs.macos-compile-admission.outputs.producer_run_attempt }} CMUX_NODE_PRODUCT_CACHE_HIT: ${{ steps.node-products.outputs.hit }} CMUX_NODE_PRODUCT_CACHE_LOOKUP_SECONDS: ${{ steps.node-products.outputs.lookup_seconds }} + CMUX_PEER_PRODUCT_HIT: ${{ steps.peer-products.outputs.hit }} + CMUX_PEER_PRODUCT_LOOKUP_SECONDS: ${{ steps.peer-products.outputs.lookup_seconds }} + CMUX_PEER_PRODUCT_TRANSFER_SECONDS: ${{ steps.peer-products.outputs.transfer_seconds }} + CMUX_PEER_PRODUCT_BYTES: ${{ steps.peer-products.outputs.bytes_transferred }} + CMUX_R2_PRODUCT_HIT: ${{ steps.r2-products.outputs.hit }} run: scripts/ci/restore-app-host-test-product.sh - name: Finalize node-local compiled product cache @@ -930,9 +964,10 @@ jobs: CMUX_PRODUCT_CONTRACT: ${{ needs.macos-compile-admission.outputs.product_contract }} CMUX_PRODUCT_SOURCE_REVISION: ${{ needs.macos-compile-admission.outputs.source_revision }} CMUX_PRODUCT_PRODUCER_RUN_ID: ${{ needs.macos-compile-admission.outputs.producer_run_id }} + CMUX_PRODUCT_PRODUCER_RUN_ATTEMPT: ${{ needs.macos-compile-admission.outputs.producer_run_attempt }} CMUX_NODE_PRODUCT_CACHE_TOKEN: ${{ steps.node-products.outputs.token }} CMUX_NODE_PRODUCT_CACHE_LEASE: ${{ steps.node-products.outputs.lease }} - CMUX_NODE_PRODUCT_SOURCE_CLASS: ${{ steps.r2-products.outputs.hit == 'true' && 'r2' || 'github' }} + CMUX_NODE_PRODUCT_SOURCE_CLASS: ${{ steps.peer-products.outputs.hit == 'true' && 'peer' || (steps.r2-products.outputs.hit == 'true' && 'r2' || 'github') }} CMUX_PRODUCT_RESTORE_SUCCEEDED: ${{ steps.restore-products.outcome == 'success' }} run: python3 scripts/ci/node_product_cache.py finalize "$RUNNER_TEMP/app-host-products/app-host-products.tar.gz" @@ -2558,6 +2593,8 @@ jobs: CMUX_NODE_PRODUCT_CACHE_ROOT: ${{ vars.CMUX_NODE_PRODUCT_CACHE_ROOT }} CMUX_NODE_PRODUCT_CACHE_MAX_BYTES: ${{ vars.CMUX_NODE_PRODUCT_CACHE_MAX_BYTES }} CMUX_NODE_PRODUCT_CACHE_WAIT_SECONDS: ${{ vars.CMUX_NODE_PRODUCT_CACHE_WAIT_SECONDS }} + CMUX_ARTIFACT_PEER_URLS: ${{ vars.CMUX_ARTIFACT_PEER_URLS }} + CMUX_ARTIFACT_PEER_TOKEN_FILE: ${{ vars.CMUX_ARTIFACT_PEER_TOKEN_FILE }} CMUX_CI_XCODE_APP: ${{ vars.CMUX_CI_XCODE_APP_MACOS_15 }} CMUX_CI_REQUIRED_MACOS_SDK_MAJOR: "26" steps: @@ -2626,12 +2663,29 @@ jobs: CMUX_PRODUCT_CONTRACT: ${{ needs.macos-compile-admission.outputs.product_contract }} CMUX_PRODUCT_SOURCE_REVISION: ${{ needs.macos-compile-admission.outputs.source_revision }} CMUX_PRODUCT_PRODUCER_RUN_ID: ${{ needs.macos-compile-admission.outputs.producer_run_id }} - CMUX_NODE_PRODUCT_CACHE_FALLBACK_SOURCE: ${{ vars.CI_ARTIFACT_R2_URL != '' && 'r2' || 'github' }} + CMUX_PRODUCT_PRODUCER_RUN_ATTEMPT: ${{ needs.macos-compile-admission.outputs.producer_run_attempt }} + CMUX_NODE_PRODUCT_CACHE_FALLBACK_SOURCE: ${{ vars.CMUX_ARTIFACT_PEER_URLS != '' && 'peer' || (vars.CI_ARTIFACT_R2_URL != '' && 'r2' || 'github') }} run: python3 scripts/ci/node_product_cache.py acquire "$RUNNER_TEMP/app-host-products" + - name: Try trusted fleet peer artifact source + id: peer-products + if: steps.node-products.outputs.hit != 'true' + continue-on-error: true + env: + ARTIFACT_ID: ${{ needs.macos-compile-admission.outputs.artifact_id }} + ARTIFACT_PROVIDER_DIGEST: ${{ needs.macos-compile-admission.outputs.artifact_digest }} + EXPECTED_SHA256: ${{ needs.macos-compile-admission.outputs.sha256 }} + CMUX_PRODUCT_CONTRACT: ${{ needs.macos-compile-admission.outputs.product_contract }} + CMUX_PRODUCT_SOURCE_REVISION: ${{ needs.macos-compile-admission.outputs.source_revision }} + CMUX_PRODUCT_PRODUCER_RUN_ID: ${{ needs.macos-compile-admission.outputs.producer_run_id }} + CMUX_PRODUCT_PRODUCER_RUN_ATTEMPT: ${{ needs.macos-compile-admission.outputs.producer_run_attempt }} + run: python3 scripts/ci/peer_product_source.py fetch "$RUNNER_TEMP/app-host-products" + + # R2 is an optional remote broker. With no configured endpoint, a peer + # miss falls straight through to the canonical GitHub artifact. - name: Try shared R2 artifact transport id: r2-products - if: steps.node-products.outputs.hit != 'true' + if: steps.node-products.outputs.hit != 'true' && steps.peer-products.outputs.hit != 'true' && vars.CI_ARTIFACT_R2_URL != '' continue-on-error: true env: GH_TOKEN: ${{ github.token }} @@ -2641,7 +2695,7 @@ jobs: run: python3 scripts/ci/restore-r2-artifact.py - name: Download compiled app-host test product - if: steps.node-products.outputs.hit != 'true' && steps.r2-products.outputs.hit != 'true' + if: steps.node-products.outputs.hit != 'true' && steps.peer-products.outputs.hit != 'true' && steps.r2-products.outputs.hit != 'true' uses: actions/download-artifact@37930b1c2abaa49bbe596cd826c3c89aef350131 # v7.0.0 with: artifact-ids: ${{ needs.macos-compile-admission.outputs.artifact_id }} @@ -2651,9 +2705,20 @@ jobs: - name: Restore compiled app-host test product id: restore-products env: + ARTIFACT_ID: ${{ needs.macos-compile-admission.outputs.artifact_id }} + ARTIFACT_PROVIDER_DIGEST: ${{ needs.macos-compile-admission.outputs.artifact_digest }} EXPECTED_SHA256: ${{ needs.macos-compile-admission.outputs.sha256 }} + CMUX_PRODUCT_CONTRACT: ${{ needs.macos-compile-admission.outputs.product_contract }} + CMUX_PRODUCT_SOURCE_REVISION: ${{ needs.macos-compile-admission.outputs.source_revision }} + CMUX_PRODUCT_PRODUCER_RUN_ID: ${{ needs.macos-compile-admission.outputs.producer_run_id }} + CMUX_PRODUCT_PRODUCER_RUN_ATTEMPT: ${{ needs.macos-compile-admission.outputs.producer_run_attempt }} CMUX_NODE_PRODUCT_CACHE_HIT: ${{ steps.node-products.outputs.hit }} CMUX_NODE_PRODUCT_CACHE_LOOKUP_SECONDS: ${{ steps.node-products.outputs.lookup_seconds }} + CMUX_PEER_PRODUCT_HIT: ${{ steps.peer-products.outputs.hit }} + CMUX_PEER_PRODUCT_LOOKUP_SECONDS: ${{ steps.peer-products.outputs.lookup_seconds }} + CMUX_PEER_PRODUCT_TRANSFER_SECONDS: ${{ steps.peer-products.outputs.transfer_seconds }} + CMUX_PEER_PRODUCT_BYTES: ${{ steps.peer-products.outputs.bytes_transferred }} + CMUX_R2_PRODUCT_HIT: ${{ steps.r2-products.outputs.hit }} run: scripts/ci/restore-app-host-test-product.sh - name: Finalize node-local compiled product cache @@ -2667,9 +2732,10 @@ jobs: CMUX_PRODUCT_CONTRACT: ${{ needs.macos-compile-admission.outputs.product_contract }} CMUX_PRODUCT_SOURCE_REVISION: ${{ needs.macos-compile-admission.outputs.source_revision }} CMUX_PRODUCT_PRODUCER_RUN_ID: ${{ needs.macos-compile-admission.outputs.producer_run_id }} + CMUX_PRODUCT_PRODUCER_RUN_ATTEMPT: ${{ needs.macos-compile-admission.outputs.producer_run_attempt }} CMUX_NODE_PRODUCT_CACHE_TOKEN: ${{ steps.node-products.outputs.token }} CMUX_NODE_PRODUCT_CACHE_LEASE: ${{ steps.node-products.outputs.lease }} - CMUX_NODE_PRODUCT_SOURCE_CLASS: ${{ steps.r2-products.outputs.hit == 'true' && 'r2' || 'github' }} + CMUX_NODE_PRODUCT_SOURCE_CLASS: ${{ steps.peer-products.outputs.hit == 'true' && 'peer' || (steps.r2-products.outputs.hit == 'true' && 'r2' || 'github') }} CMUX_PRODUCT_RESTORE_SUCCEEDED: ${{ steps.restore-products.outcome == 'success' }} run: python3 scripts/ci/node_product_cache.py finalize "$RUNNER_TEMP/app-host-products/app-host-products.tar.gz" diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 45f29579fb33..afdd1b95ea0d 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -156,7 +156,7 @@ jobs: ci_router_changed=true ghosttykit_guard_only=false ;; - scripts/ci/build_input_fingerprint.py|scripts/ci/find_admitted_build.py|scripts/ci/app_host_test_products.py|scripts/ci/compile-app-host-test-product.sh|scripts/ci/product_input_identity.py|scripts/ci/restore-app-host-test-product.sh|scripts/ci/reuse_app_host_products.py|scripts/ci/sanitize-xcode-source-packages-cache.py|scripts/ci/persistent_mac_route.py|scripts/ci/web_validation.py) + scripts/ci/build_input_fingerprint.py|scripts/ci/find_admitted_build.py|scripts/ci/app_host_test_products.py|scripts/ci/compile-app-host-test-product.sh|scripts/ci/product_input_identity.py|scripts/ci/peer_product_source.py|scripts/ci/restore-app-host-test-product.sh|scripts/ci/reuse_app_host_products.py|scripts/ci/sanitize-xcode-source-packages-cache.py|scripts/ci/persistent_mac_route.py|scripts/ci/web_validation.py) # These have explicit owning lanes and names that cannot # shadow Python stdlib imports used by the router. ghosttykit_guard_only=false diff --git a/scripts/ci/detect_ci_change_areas.py b/scripts/ci/detect_ci_change_areas.py index 54528e9f8c18..81015ff9280a 100755 --- a/scripts/ci/detect_ci_change_areas.py +++ b/scripts/ci/detect_ci_change_areas.py @@ -72,6 +72,7 @@ def is_other_workflow_config(path: str) -> bool: "scripts/ci/app_host_test_products.py", "scripts/ci/compile-app-host-test-product.sh", "scripts/ci/product_input_identity.py", + "scripts/ci/peer_product_source.py", "scripts/ci/restore-app-host-test-product.sh", "scripts/ci/reuse_app_host_products.py", "scripts/ci/sanitize-xcode-source-packages-cache.py", diff --git a/scripts/ci/detect_linux_guard_changes.py b/scripts/ci/detect_linux_guard_changes.py index 77c1566cf844..5d381c7e649f 100644 --- a/scripts/ci/detect_linux_guard_changes.py +++ b/scripts/ci/detect_linux_guard_changes.py @@ -30,6 +30,7 @@ "scripts/ci/app_host_test_products.py", "scripts/ci/compile-app-host-test-product.sh", "scripts/ci/product_input_identity.py", + "scripts/ci/peer_product_source.py", "scripts/ci/restore-app-host-test-product.sh", "scripts/ci/reuse_app_host_products.py", "scripts/ci/sanitize-xcode-source-packages-cache.py", diff --git a/scripts/ci/node_product_cache.py b/scripts/ci/node_product_cache.py index c7cbea8fba4d..e4dc2d802382 100644 --- a/scripts/ci/node_product_cache.py +++ b/scripts/ci/node_product_cache.py @@ -27,7 +27,7 @@ from pathlib import Path from typing import Callable, Iterator -SCHEMA_GENERATION = 1 +SCHEMA_GENERATION = 2 FORMAT_GENERATION = "app-host-products-tar-gz-v1" PROVIDER = "github-actions" ARCHIVE_NAME = "app-host-products.tar.gz" @@ -68,6 +68,7 @@ class Identity: product_contract: str source_revision: str producer_run_id: int + producer_run_attempt: int = 1 schema_generation: int = SCHEMA_GENERATION format_generation: str = FORMAT_GENERATION provider: str = PROVIDER @@ -84,6 +85,7 @@ def as_dict(self) -> dict: "product_contract": self.product_contract, "source_revision": self.source_revision, "producer_run_id": self.producer_run_id, + "producer_run_attempt": self.producer_run_attempt, } def key(self) -> str: @@ -109,6 +111,9 @@ def from_env(cls, env=os.environ) -> "Identity": producer_run_id=_positive_int( env.get("CMUX_PRODUCT_PRODUCER_RUN_ID", ""), "producer run id" ), + producer_run_attempt=_positive_int( + env.get("CMUX_PRODUCT_PRODUCER_RUN_ATTEMPT", "1"), "producer run attempt" + ), ) @@ -262,7 +267,7 @@ def _metadata_matches(metadata: dict, identity: Identity) -> bool: and metadata.get("object_digest") == identity.archive_digest and isinstance(metadata.get("size"), int) and metadata["size"] > 0 - and metadata.get("source_class") in {"github", "r2", "producer-local"} + and metadata.get("source_class") in {"github", "r2", "peer", "producer-local"} ) @@ -386,6 +391,7 @@ def _stats_update(store: Store, **increments) -> dict: "evictions": 0, "evicted_bytes": 0, "bytes_avoided_github": 0, + "bytes_avoided_peer": 0, "bytes_avoided_r2": 0, } for field, amount in increments.items(): @@ -426,6 +432,7 @@ def _snapshot(store: Store, stats: dict | None = None) -> dict: "evictions": evictions, "eviction_rate": round(evictions / lookups, 4) if lookups else 0.0, "bytes_avoided_github": int(stats.get("bytes_avoided_github", 0)), + "bytes_avoided_peer": int(stats.get("bytes_avoided_peer", 0)), "bytes_avoided_r2": int(stats.get("bytes_avoided_r2", 0)), } @@ -537,7 +544,9 @@ def _hit_locked( avoided = metadata["size"] fallback = os.environ.get("CMUX_NODE_PRODUCT_CACHE_FALLBACK_SOURCE", "").strip() increments = {"hits": 1} - if fallback == "r2": + if fallback == "peer": + increments["bytes_avoided_peer"] = avoided + elif fallback == "r2": increments["bytes_avoided_r2"] = avoided elif fallback == "github": increments["bytes_avoided_github"] = avoided @@ -671,6 +680,28 @@ def github_metadata(identity: Identity) -> dict: return json.loads(raw) +def same_run_provider_metadata(identity: Identity) -> dict: + """Reconstruct provider metadata only for this exact producing workflow attempt. + + Same-run GitHub workflow outputs establish the accepted producer identity. + This keeps a verified peer hit usable during a GitHub API outage without + making the peer a new trust root. + """ + if ( + os.environ.get("GITHUB_REPOSITORY", "").casefold() != identity.repository.casefold() + or os.environ.get("GITHUB_RUN_ID", "") != str(identity.producer_run_id) + or os.environ.get("GITHUB_RUN_ATTEMPT", "1") != str(identity.producer_run_attempt) + ): + raise ValueError("peer product producer is not this workflow attempt") + return { + "id": identity.artifact_id, + "expired": False, + "digest": "sha256:" + identity.provider_digest, + "workflow_run": {"id": identity.producer_run_id}, + "created_at": None, + } + + def _verify_provider(identity: Identity, metadata: dict) -> str | None: digest = metadata.get("digest", "") run = metadata.get("workflow_run") @@ -824,7 +855,7 @@ def finalize( ) -> dict: if store is None: return {"status": "disabled"} - if source_class not in {"github", "r2", "producer-local"}: + if source_class not in {"github", "r2", "peer", "producer-local"}: source_class = "github" key = identity.key() if not restore_succeeded: @@ -1079,15 +1110,19 @@ def main() -> None: elif command == "finalize": if len(sys.argv) != 3: raise SystemExit("usage: node-product-cache.py finalize ARCHIVE") + source_class = os.environ.get("CMUX_NODE_PRODUCT_SOURCE_CLASS", "github") result = finalize( store, identity, Path(sys.argv[2]), token=os.environ.get("CMUX_NODE_PRODUCT_CACHE_TOKEN", ""), lease_token=os.environ.get("CMUX_NODE_PRODUCT_CACHE_LEASE", ""), - source_class=os.environ.get("CMUX_NODE_PRODUCT_SOURCE_CLASS", "github"), + source_class=source_class, restore_succeeded=os.environ.get("CMUX_PRODUCT_RESTORE_SUCCEEDED") == "true", budget=budget_bytes(), + provider_metadata=( + same_run_provider_metadata if source_class == "peer" else github_metadata + ), ) else: if len(sys.argv) != 3: diff --git a/scripts/ci/peer_product_source.py b/scripts/ci/peer_product_source.py new file mode 100644 index 000000000000..0f36a55a23b7 --- /dev/null +++ b/scripts/ci/peer_product_source.py @@ -0,0 +1,720 @@ +#!/usr/bin/env python3 +"""Trusted exact-object peer source for immutable compiled CI products. + +The peer service exposes only HEAD/GET for one caller-supplied immutable object +key. It never lists cache contents, exposes cache paths, or accepts writes. +Client requests use one monotonic deadline across underlying receives, and the +server bounds client socket time plus concurrent requests so a slow peer falls +through to the existing sources. Transport is acceleration only: callers still +verify the content digest and run the canonical product restore validator before +publishing locally. +""" +from __future__ import annotations + +import argparse +import contextlib +import hashlib +import hmac +import http.client +import json +import os +import re +import socket +import ssl +import stat +import tempfile +import threading +import time +from dataclasses import dataclass +from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer +from pathlib import Path +from typing import Callable, Iterator +from urllib.parse import urlsplit + +import node_product_cache as cache + +MAX_PEERS = 8 +MAX_OBJECT_BYTES = 20 * 1024**3 +DEFAULT_LOOKUP_TIMEOUT_SECONDS = 5.0 +DEFAULT_TRANSFER_TIMEOUT_SECONDS = 180.0 +DEFAULT_SERVER_CLIENT_TIMEOUT_SECONDS = 30.0 +DEFAULT_SERVER_MAX_ACTIVE_REQUESTS = 16 +OBJECT_PATH_PREFIX = "/v1/objects/" +OBJECT_KEY_RE = re.compile(r"[a-f0-9]{64}") + + +class PeerUnavailable(RuntimeError): + """One peer cannot serve the exact requested object.""" + + +@dataclass(frozen=True) +class PeerSource: + url: str + + def __post_init__(self) -> None: + parsed = urlsplit(self.url) + if ( + parsed.scheme != "https" + or not parsed.netloc + or parsed.username + or parsed.password + or parsed.query + or parsed.fragment + or parsed.path not in ("", "/") + ): + raise ValueError("peer source must be one HTTPS origin") + + @property + def origin(self) -> str: + return self.url.rstrip("/") + + +@dataclass(frozen=True) +class PeerAvailability: + object_key: str + schema_generation: int + size_bytes: int + content_digest: str + + def validate_for(self, identity: cache.Identity) -> None: + if ( + self.object_key != identity.key() + or self.schema_generation != cache.SCHEMA_GENERATION + or self.content_digest != identity.archive_digest + or not 0 < self.size_bytes <= MAX_OBJECT_BYTES + ): + raise PeerUnavailable("peer offer does not match exact object identity") + + +@dataclass +class OpenLocalObject: + path: Path + size_bytes: int + content_digest: str + identity: cache.Identity + + +def _identity_from_metadata(value: dict) -> cache.Identity: + raw = value.get("identity") + if not isinstance(raw, dict): + raise PeerUnavailable("peer object metadata has no identity") + try: + return cache.Identity(**raw) + except (TypeError, ValueError) as error: + raise PeerUnavailable("peer object identity is invalid") from error + + +def local_availability(store: cache.Store, object_key: str) -> PeerAvailability | None: + if not OBJECT_KEY_RE.fullmatch(object_key): + return None + with store.lock(object_key): + metadata = cache._read_json(store.entry(object_key) / cache.METADATA_NAME) + if metadata is None: + return None + try: + identity = _identity_from_metadata(metadata) + except PeerUnavailable: + return None + if identity.key() != object_key: + return None + validated = cache._validate_entry_locked(store, identity) + if validated is None: + return None + metadata, _ = validated + return PeerAvailability( + object_key=object_key, + schema_generation=int(metadata["schema_generation"]), + size_bytes=int(metadata["size"]), + content_digest=str(metadata["object_digest"]), + ) + + +@contextlib.contextmanager +def open_local_object( + store: cache.Store, + object_key: str, + *, + draining: Callable[[], bool], +) -> Iterator[OpenLocalObject]: + if not OBJECT_KEY_RE.fullmatch(object_key) or draining(): + raise PeerUnavailable("peer source is unavailable") + lease = "" + identity = None + obj = None + metadata = None + with store.lock(object_key): + if draining(): + raise PeerUnavailable("peer source is draining") + raw = cache._read_json(store.entry(object_key) / cache.METADATA_NAME) + if raw is None: + raise PeerUnavailable("peer object is absent") + identity = _identity_from_metadata(raw) + if identity.key() != object_key: + raise PeerUnavailable("peer object identity mismatch") + validated = cache._validate_entry_locked(store, identity) + if validated is None: + raise PeerUnavailable("peer object is unavailable") + metadata, obj = validated + lease = cache._create_lease_locked(store, object_key) + try: + assert identity is not None and obj is not None and metadata is not None + yield OpenLocalObject( + path=obj, + size_bytes=int(metadata["size"]), + content_digest=str(metadata["object_digest"]), + identity=identity, + ) + finally: + if lease: + with contextlib.suppress(OSError): + with store.lock(object_key): + cache._release_lease_locked(store, object_key, lease) + + +def _read_secret(path: Path) -> str: + if not path.is_absolute(): + raise ValueError("peer token file must be absolute") + info = path.lstat() + if not stat.S_ISREG(info.st_mode) or stat.S_ISLNK(info.st_mode): + raise ValueError("peer token must be a regular file") + if info.st_uid != os.geteuid() or info.st_mode & 0o077: + raise ValueError("peer token file must be private and owned by this user") + if info.st_size <= 0 or info.st_size > 4096: + raise ValueError("peer token size is invalid") + token = path.read_text().strip() + if len(token) < 32 or any(ord(char) < 0x21 or ord(char) > 0x7E for char in token): + raise ValueError("peer token is invalid") + return token + + +def configured_sources(env=os.environ) -> list[PeerSource]: + raw = env.get("CMUX_ARTIFACT_PEER_URLS", "").strip() + if not raw: + return [] + values = [value.strip() for value in raw.split(",") if value.strip()] + if len(values) > MAX_PEERS: + return [] + try: + return [PeerSource(value) for value in values] + except ValueError: + return [] + + +def configured_token_loader(env=os.environ) -> Callable[[PeerSource], str]: + raw = env.get("CMUX_ARTIFACT_PEER_TOKEN_FILE", "").strip() + if not raw: + def unavailable(_source: PeerSource) -> str: + raise PeerUnavailable("peer token file is unavailable") + return unavailable + path = Path(raw).expanduser() + + def load(_source: PeerSource) -> str: + try: + return _read_secret(path) + except (OSError, ValueError) as error: + raise PeerUnavailable("peer token is unavailable") from error + + return load + + +def _connection(source: PeerSource, timeout: float) -> tuple[http.client.HTTPSConnection, str]: + parsed = urlsplit(source.origin) + host = parsed.hostname + if host is None: + raise PeerUnavailable("peer host is invalid") + port = parsed.port or 443 + return ( + http.client.HTTPSConnection( + host, + port, + timeout=max(0.001, timeout), + context=ssl.create_default_context(), + ), + parsed.netloc, + ) + + +def _request_deadline(timeout: float) -> float: + """Return one absolute deadline for the full peer HTTP request.""" + return time.monotonic() + max(0.1, timeout) + + +def _remaining_timeout(deadline: float) -> float: + """Return time left before deadline or fail the peer request.""" + remaining = deadline - time.monotonic() + if remaining <= 0: + raise PeerUnavailable("peer request deadline exceeded") + return max(0.001, remaining) + + +def _arm_connection_deadline( + connection: http.client.HTTPSConnection, + deadline: float, +) -> None: + """Apply the remaining absolute deadline to connect/read/write socket work.""" + remaining = _remaining_timeout(deadline) + connection.timeout = remaining + if connection.sock is not None: + connection.sock.settimeout(remaining) + + +def _read_response_once( + connection: http.client.HTTPSConnection, + response: http.client.HTTPResponse, + deadline: float, + size: int, +) -> bytes: + """Read with a fresh remaining timeout for at most one buffered raw receive.""" + _arm_connection_deadline(connection, deadline) + return response.read1(size) + + +def probe_http( + source: PeerSource, + object_key: str, + token: str, + *, + timeout: float = DEFAULT_LOOKUP_TIMEOUT_SECONDS, +) -> PeerAvailability | None: + if not OBJECT_KEY_RE.fullmatch(object_key): + return None + deadline = _request_deadline(timeout) + connection, authority = _connection(source, _remaining_timeout(deadline)) + try: + _arm_connection_deadline(connection, deadline) + connection.request( + "HEAD", + OBJECT_PATH_PREFIX + object_key, + headers={ + "Authorization": f"Bearer {token}", + "Host": authority, + "Accept": "application/octet-stream", + }, + ) + _arm_connection_deadline(connection, deadline) + response = connection.getresponse() + while True: + if not _read_response_once(connection, response, deadline, 64 * 1024): + break + if response.status == 404: + return None + if response.status != 200: + raise PeerUnavailable(f"peer probe returned HTTP {response.status}") + try: + return PeerAvailability( + object_key=response.getheader("X-Cmux-Object-Key", ""), + schema_generation=int(response.getheader("X-Cmux-Object-Schema", "0")), + size_bytes=int(response.getheader("Content-Length", "0")), + content_digest=response.getheader("X-Cmux-Content-Sha256", ""), + ) + except (TypeError, ValueError) as error: + raise PeerUnavailable("peer availability response is invalid") from error + except (OSError, ssl.SSLError, http.client.HTTPException) as error: + raise PeerUnavailable("peer probe failed") from error + finally: + connection.close() + + +def transfer_http( + source: PeerSource, + object_key: str, + token: str, + target: Path, + size: int, + *, + timeout: float = DEFAULT_TRANSFER_TIMEOUT_SECONDS, +) -> None: + if not OBJECT_KEY_RE.fullmatch(object_key) or not 0 < size <= MAX_OBJECT_BYTES: + raise PeerUnavailable("peer transfer request is invalid") + deadline = _request_deadline(timeout) + connection, authority = _connection(source, _remaining_timeout(deadline)) + copied = 0 + try: + _arm_connection_deadline(connection, deadline) + connection.request( + "GET", + OBJECT_PATH_PREFIX + object_key, + headers={ + "Authorization": f"Bearer {token}", + "Host": authority, + "Accept": "application/octet-stream", + }, + ) + _arm_connection_deadline(connection, deadline) + response = connection.getresponse() + if response.status != 200: + while True: + if not _read_response_once(connection, response, deadline, 64 * 1024): + break + raise PeerUnavailable(f"peer fetch returned HTTP {response.status}") + try: + declared = int(response.getheader("Content-Length", "0")) + except ValueError as error: + raise PeerUnavailable("peer fetch size is invalid") from error + if declared != size: + raise PeerUnavailable("peer fetch size changed after probe") + with target.open("xb") as output: + while True: + chunk = _read_response_once( + connection, + response, + deadline, + min(1024 * 1024, size - copied + 1), + ) + if not chunk: + break + copied += len(chunk) + if copied > size: + raise PeerUnavailable("peer sent too many bytes") + output.write(chunk) + output.flush() + os.fsync(output.fileno()) + if copied != size: + raise PeerUnavailable("peer transfer ended early") + except (OSError, ssl.SSLError, http.client.HTTPException) as error: + raise PeerUnavailable("peer transfer failed") from error + finally: + connection.close() + + +def _sha256(path: Path) -> str: + digest = hashlib.sha256() + with path.open("rb") as source: + for chunk in iter(lambda: source.read(1024 * 1024), b""): + digest.update(chunk) + return digest.hexdigest() + + +def fetch_exact( + identity: cache.Identity, + destination: Path, + sources: list[PeerSource], + *, + probe: Callable[[PeerSource, str, str], PeerAvailability | None] = probe_http, + transfer: Callable[[PeerSource, str, str, Path, int], None] = transfer_http, + token_loader: Callable[[PeerSource], str] | None = None, +) -> dict: + started = time.monotonic() + object_key = identity.key() + if destination.exists() or not sources: + return { + "status": "miss", + "hit": False, + "source": "", + "source_index": -1, + "lookup_seconds": round(time.monotonic() - started, 6), + "transfer_seconds": 0.0, + "bytes_transferred": 0, + } + token_loader = token_loader or configured_token_loader() + destination.parent.mkdir(parents=True, exist_ok=True) + for index, source in enumerate(sources[:MAX_PEERS]): + lookup_started = time.monotonic() + try: + token = token_loader(source) + offer = probe(source, object_key, token) + lookup_seconds = time.monotonic() - lookup_started + if offer is None: + continue + offer.validate_for(identity) + transfer_started = time.monotonic() + with tempfile.TemporaryDirectory(prefix="cmux-peer-product-", dir=destination.parent) as raw: + staging = Path(raw) + obj = staging / cache.ARCHIVE_NAME + transfer(source, object_key, token, obj, offer.size_bytes) + transfer_seconds = time.monotonic() - transfer_started + if ( + obj.stat().st_size != offer.size_bytes + or _sha256(obj) != identity.archive_digest + ): + raise PeerUnavailable("peer object digest mismatch") + products = staging / "products" + products.mkdir() + os.rename(obj, products / cache.ARCHIVE_NAME) + if destination.exists(): + raise PeerUnavailable("peer destination became occupied") + os.rename(products, destination) + return { + "status": "hit", + "hit": True, + "source": "peer", + "source_index": index, + "object_key": object_key, + "archive_bytes": offer.size_bytes, + "lookup_seconds": round(lookup_seconds, 6), + "transfer_seconds": round(transfer_seconds, 6), + "bytes_transferred": offer.size_bytes, + } + except PeerUnavailable: + continue + except (OSError, ValueError, TypeError): + continue + return { + "status": "miss", + "hit": False, + "source": "", + "source_index": -1, + "object_key": object_key, + "lookup_seconds": round(time.monotonic() - started, 6), + "transfer_seconds": 0.0, + "bytes_transferred": 0, + } + + +def _append_outputs(values: dict) -> None: + output_path = os.environ.get("GITHUB_OUTPUT") + if not output_path: + return + with open(output_path, "a") as output: + for key, value in values.items(): + if isinstance(value, bool): + value = str(value).lower() + output.write(f"{key}={value}\n") + + +def _report(result: dict) -> None: + record = { + "run_id": os.environ.get("GITHUB_RUN_ID"), + "job": os.environ.get("GITHUB_JOB"), + "runner_name": os.environ.get("RUNNER_NAME"), + **result, + } + print("CMUX_PEER_PRODUCT_SOURCE " + json.dumps(record, sort_keys=True)) + summary = os.environ.get("GITHUB_STEP_SUMMARY") + if summary: + with contextlib.suppress(OSError): + with open(summary, "a") as handle: + handle.write("### Trusted peer compiled product source\n\n```json\n") + handle.write(json.dumps(record, indent=2, sort_keys=True)) + handle.write("\n```\n") + + +def _draining(marker: Path | None) -> bool: + return marker is not None and marker.exists() + + +class PeerRequestHandler(BaseHTTPRequestHandler): + server_version = "cmux-peer-artifact/1" + protocol_version = "HTTP/1.1" + + def log_message(self, _format: str, *_args) -> None: + return + + def _authorized(self) -> bool: + supplied = self.headers.get("Authorization", "") + expected = f"Bearer {self.server.peer_token}" + return hmac.compare_digest(supplied, expected) + + def _object_key(self) -> str | None: + if not self.path.startswith(OBJECT_PATH_PREFIX): + return None + value = self.path[len(OBJECT_PATH_PREFIX):] + return value if OBJECT_KEY_RE.fullmatch(value) else None + + def _send_headers(self, offer: PeerAvailability) -> None: + self.send_response(200) + self.send_header("Content-Type", "application/octet-stream") + self.send_header("Content-Length", str(offer.size_bytes)) + self.send_header("X-Cmux-Object-Key", offer.object_key) + self.send_header("X-Cmux-Object-Schema", str(offer.schema_generation)) + self.send_header("X-Cmux-Content-Sha256", offer.content_digest) + self.send_header("Cache-Control", "private, immutable") + self.end_headers() + + def _refuse(self, status: int) -> None: + self.send_response(status) + self.send_header("Content-Length", "0") + self.end_headers() + + def do_HEAD(self) -> None: + if not self._authorized(): + self._refuse(401) + return + key = self._object_key() + if key is None: + self._refuse(404) + return + if self.server.draining(): + self._refuse(503) + return + offer = local_availability(self.server.store, key) + if offer is None: + self._refuse(404) + return + self._send_headers(offer) + + def do_GET(self) -> None: + if not self._authorized(): + self._refuse(401) + return + key = self._object_key() + if key is None: + self._refuse(404) + return + try: + with open_local_object( + self.server.store, + key, + draining=self.server.draining, + ) as opened: + offer = PeerAvailability( + object_key=key, + schema_generation=cache.SCHEMA_GENERATION, + size_bytes=opened.size_bytes, + content_digest=opened.content_digest, + ) + self._send_headers(offer) + with opened.path.open("rb") as source: + while chunk := source.read(1024 * 1024): + self.wfile.write(chunk) + except PeerUnavailable: + self._refuse(503) + except (BrokenPipeError, ConnectionResetError, socket.timeout): + return + + def do_POST(self) -> None: + self._refuse(405) + + do_PUT = do_POST + do_DELETE = do_POST + do_PATCH = do_POST + + +class PeerHTTPServer(ThreadingHTTPServer): + daemon_threads = True + + def __init__( + self, + address, + handler, + *, + store: cache.Store, + token: str, + drain_marker: Path | None, + client_timeout_seconds: float = DEFAULT_SERVER_CLIENT_TIMEOUT_SECONDS, + max_active_requests: int = DEFAULT_SERVER_MAX_ACTIVE_REQUESTS, + ) -> None: + if client_timeout_seconds <= 0: + raise ValueError("peer client timeout must be positive") + if max_active_requests <= 0: + raise ValueError("peer active request limit must be positive") + super().__init__(address, handler) + self.store = store + self.peer_token = token + self._drain_marker = drain_marker + self.client_timeout_seconds = client_timeout_seconds + self.max_active_requests = max_active_requests + self._request_slots = threading.BoundedSemaphore(max_active_requests) + + def get_request(self): + request, client_address = super().get_request() + try: + # The listener defers TLS handshakes. Keep the accept loop free of + # client-controlled handshake work; the bounded request worker will + # perform it lazily under this socket deadline. + request.settimeout(self.client_timeout_seconds) + return request, client_address + except BaseException: + request.close() + raise + + def process_request(self, request, client_address) -> None: + if not self._request_slots.acquire(blocking=False): + self.shutdown_request(request) + return + try: + super().process_request(request, client_address) + except BaseException: + self._request_slots.release() + raise + + def process_request_thread(self, request, client_address) -> None: + try: + super().process_request_thread(request, client_address) + finally: + self._request_slots.release() + + def draining(self) -> bool: + return _draining(self._drain_marker) + + +def serve(args: argparse.Namespace) -> None: + store = cache.Store(args.root.resolve()) + token = _read_secret(args.token_file.resolve()) + server = PeerHTTPServer( + (args.bind, args.port), + PeerRequestHandler, + store=store, + token=token, + drain_marker=args.drain_marker.resolve() if args.drain_marker else None, + client_timeout_seconds=args.client_timeout_seconds, + max_active_requests=args.max_active_requests, + ) + context = ssl.SSLContext(ssl.PROTOCOL_TLS_SERVER) + context.minimum_version = ssl.TLSVersion.TLSv1_2 + context.load_cert_chain(args.cert.resolve(), args.key.resolve()) + server.socket = context.wrap_socket( + server.socket, + server_side=True, + do_handshake_on_connect=False, + ) + try: + server.serve_forever() + finally: + server.server_close() + + +def main() -> None: + parser = argparse.ArgumentParser(description="CMUX exact immutable peer artifact source") + sub = parser.add_subparsers(dest="command", required=True) + fetch_parser = sub.add_parser("fetch") + fetch_parser.add_argument("destination", type=Path) + serve_parser = sub.add_parser("serve") + serve_parser.add_argument("--root", type=Path, required=True) + serve_parser.add_argument("--bind", default="127.0.0.1") + serve_parser.add_argument("--port", type=int, default=9443) + serve_parser.add_argument("--cert", type=Path, required=True) + serve_parser.add_argument("--key", type=Path, required=True) + serve_parser.add_argument("--token-file", type=Path, required=True) + serve_parser.add_argument("--drain-marker", type=Path) + serve_parser.add_argument( + "--client-timeout-seconds", + type=float, + default=DEFAULT_SERVER_CLIENT_TIMEOUT_SECONDS, + ) + serve_parser.add_argument( + "--max-active-requests", + type=int, + default=DEFAULT_SERVER_MAX_ACTIVE_REQUESTS, + ) + args = parser.parse_args() + + if args.command == "serve": + serve(args) + return + + result = { + "status": "miss", + "hit": False, + "source": "", + "source_index": -1, + "lookup_seconds": 0.0, + "transfer_seconds": 0.0, + "bytes_transferred": 0, + } + try: + identity = cache.Identity.from_env() + result = fetch_exact( + identity, + args.destination, + configured_sources(), + token_loader=configured_token_loader(), + ) + except (OSError, ValueError, TypeError, PeerUnavailable): + pass + _append_outputs(result) + _report(result) + + +if __name__ == "__main__": + main() diff --git a/scripts/ci/restore-app-host-test-product.sh b/scripts/ci/restore-app-host-test-product.sh index 91987eda8102..326716d8b5ba 100755 --- a/scripts/ci/restore-app-host-test-product.sh +++ b/scripts/ci/restore-app-host-test-product.sh @@ -30,11 +30,27 @@ import time archive = Path(os.environ["RUNNER_TEMP"]) / "app-host-products/app-host-products.tar.gz" elapsed = max(0.0, (time.monotonic_ns() - int(os.environ["CMUX_RESTORE_STARTED_NS"])) / 1_000_000_000) +local_hit = os.environ.get("CMUX_NODE_PRODUCT_CACHE_HIT") == "true" +peer_hit = os.environ.get("CMUX_PEER_PRODUCT_HIT") == "true" +r2_hit = os.environ.get("CMUX_R2_PRODUCT_HIT") == "true" record = { + "repository": os.environ["GITHUB_REPOSITORY"], + "artifact_id": int(os.environ["ARTIFACT_ID"]), + "provider_digest": os.environ["ARTIFACT_PROVIDER_DIGEST"], + "archive_sha256": os.environ["EXPECTED_SHA256"], + "product_contract": os.environ["CMUX_PRODUCT_CONTRACT"], + "source_revision": os.environ["CMUX_PRODUCT_SOURCE_REVISION"], + "producer_run_id": int(os.environ["CMUX_PRODUCT_PRODUCER_RUN_ID"]), + "producer_run_attempt": int(os.environ["CMUX_PRODUCT_PRODUCER_RUN_ATTEMPT"]), "archive_bytes": archive.stat().st_size, "elapsed_seconds": round(elapsed, 6), - "local_hit": os.environ.get("CMUX_NODE_PRODUCT_CACHE_HIT") == "true", + "lookup_source": "local" if local_hit else "peer" if peer_hit else "r2" if r2_hit else "github", + "local_hit": local_hit, "lookup_seconds": float(os.environ.get("CMUX_NODE_PRODUCT_CACHE_LOOKUP_SECONDS") or 0), + "peer_hit": peer_hit, + "peer_lookup_seconds": float(os.environ.get("CMUX_PEER_PRODUCT_LOOKUP_SECONDS") or 0), + "peer_transfer_seconds": float(os.environ.get("CMUX_PEER_PRODUCT_TRANSFER_SECONDS") or 0), + "peer_bytes_transferred": int(os.environ.get("CMUX_PEER_PRODUCT_BYTES") or 0), "run_id": os.environ.get("GITHUB_RUN_ID"), "job": os.environ.get("GITHUB_JOB"), "shard": os.environ.get("CMUX_APP_HOST_SHARD"), diff --git a/tests/test_ci_change_areas.py b/tests/test_ci_change_areas.py index 14e0a9941996..d7f74a5b5ac2 100644 --- a/tests/test_ci_change_areas.py +++ b/tests/test_ci_change_areas.py @@ -274,6 +274,7 @@ def test_macos_test_product_ci_helpers_run_admission_without_web_or_release() -> "scripts/ci/app_host_test_products.py", "scripts/ci/compile-app-host-test-product.sh", "scripts/ci/product_input_identity.py", + "scripts/ci/peer_product_source.py", "scripts/ci/restore-app-host-test-product.sh", "scripts/ci/reuse_app_host_products.py", "scripts/ci/sanitize-xcode-source-packages-cache.py", @@ -2220,6 +2221,63 @@ def test_compiled_product_cache_is_opt_in_on_persistent_macos_lanes() -> None: assert "CMUX_NODE_PRODUCT_CACHE_ROOT: ${{ vars.CMUX_NODE_PRODUCT_CACHE_ROOT }}" in block assert "CMUX_NODE_PRODUCT_CACHE_MAX_BYTES: ${{ vars.CMUX_NODE_PRODUCT_CACHE_MAX_BYTES }}" in block assert "CMUX_NODE_PRODUCT_CACHE_WAIT_SECONDS: ${{ vars.CMUX_NODE_PRODUCT_CACHE_WAIT_SECONDS }}" in block + assert "CMUX_ARTIFACT_PEER_URLS: ${{ vars.CMUX_ARTIFACT_PEER_URLS }}" in block + assert "CMUX_ARTIFACT_PEER_TOKEN_FILE: ${{ vars.CMUX_ARTIFACT_PEER_TOKEN_FILE }}" in block + + +def test_product_restore_receipt_binds_immutable_product_identity() -> None: + required_env = ( + "ARTIFACT_ID: ${{ needs.macos-compile-admission.outputs.artifact_id }}", + "ARTIFACT_PROVIDER_DIGEST: ${{ needs.macos-compile-admission.outputs.artifact_digest }}", + "EXPECTED_SHA256: ${{ needs.macos-compile-admission.outputs.sha256 }}", + "CMUX_PRODUCT_CONTRACT: ${{ needs.macos-compile-admission.outputs.product_contract }}", + "CMUX_PRODUCT_SOURCE_REVISION: ${{ needs.macos-compile-admission.outputs.source_revision }}", + "CMUX_PRODUCT_PRODUCER_RUN_ID: ${{ needs.macos-compile-admission.outputs.producer_run_id }}", + "CMUX_PRODUCT_PRODUCER_RUN_ATTEMPT: ${{ needs.macos-compile-admission.outputs.producer_run_attempt }}", + ) + for job_name in ("app-host-unit-tests", "tests-build-and-lag"): + block = workflow_job_block(job_name, MACOS_WORKFLOW) + restore = block[block.index(" - name: Restore compiled app-host test product"):] + restore = restore[:restore.index("\n - name:", 1)] + for binding in required_env: + assert binding in restore, (job_name, binding) + + script = (ROOT / "scripts/ci/restore-app-host-test-product.sh").read_text(encoding="utf-8") + for field in ( + '"repository": os.environ["GITHUB_REPOSITORY"]', + '"artifact_id": int(os.environ["ARTIFACT_ID"])', + '"provider_digest": os.environ["ARTIFACT_PROVIDER_DIGEST"]', + '"archive_sha256": os.environ["EXPECTED_SHA256"]', + '"product_contract": os.environ["CMUX_PRODUCT_CONTRACT"]', + '"source_revision": os.environ["CMUX_PRODUCT_SOURCE_REVISION"]', + '"producer_run_id": int(os.environ["CMUX_PRODUCT_PRODUCER_RUN_ID"])', + '"producer_run_attempt": int(os.environ["CMUX_PRODUCT_PRODUCER_RUN_ATTEMPT"])', + ): + assert field in script + + +def test_compiled_product_source_order_is_local_peer_r2_github() -> None: + for job_name in ("app-host-unit-tests", "tests-build-and-lag"): + block = workflow_job_block(job_name, MACOS_WORKFLOW) + assert block.index("Try node-local compiled product cache") < block.index("Try trusted fleet peer artifact source") + assert block.index("Try trusted fleet peer artifact source") < block.index("Try shared R2 artifact transport") + assert block.index("Try shared R2 artifact transport") < block.index("Download compiled app-host test product") + + +def test_r2_transport_is_an_explicit_optional_remote_broker() -> None: + expected_condition = ( + "if: steps.node-products.outputs.hit != 'true' && " + "steps.peer-products.outputs.hit != 'true' && " + "vars.CI_ARTIFACT_R2_URL != ''" + ) + for job_name in ("app-host-unit-tests", "tests-build-and-lag"): + block = workflow_job_block(job_name, MACOS_WORKFLOW) + start = block.index(" - name: Try shared R2 artifact transport") + step = block[start:] + next_step = step.index("\n - name:", 1) + r2_step = step[:next_step] + assert expected_condition in r2_step, job_name + assert "CI_ARTIFACT_R2_URL: ${{ vars.CI_ARTIFACT_R2_URL }}" in r2_step def test_macos_jobs_use_lane_specific_xcode_pin_vars() -> None: diff --git a/tests/test_ci_linux_guard_routing.py b/tests/test_ci_linux_guard_routing.py index 55680ef191ab..fa34d3b66eed 100644 --- a/tests/test_ci_linux_guard_routing.py +++ b/tests/test_ci_linux_guard_routing.py @@ -174,6 +174,7 @@ def test_macos_admission_helpers_run_only_workflow_guard_contracts(self): "scripts/ci/app_host_test_products.py", "scripts/ci/compile-app-host-test-product.sh", "scripts/ci/product_input_identity.py", + "scripts/ci/peer_product_source.py", "scripts/ci/restore-app-host-test-product.sh", "scripts/ci/reuse_app_host_products.py", "scripts/ci/sanitize-xcode-source-packages-cache.py", diff --git a/tests/test_node_product_cache.py b/tests/test_node_product_cache.py index 2b9b84fb162f..dfedeb951190 100644 --- a/tests/test_node_product_cache.py +++ b/tests/test_node_product_cache.py @@ -23,6 +23,12 @@ sys.modules[spec.name] = cache spec.loader.exec_module(cache) +PEER_MODULE = ROOT / "scripts/ci/peer_product_source.py" +peer_spec = importlib.util.spec_from_file_location("peer_product_source", PEER_MODULE) +peer = importlib.util.module_from_spec(peer_spec) +sys.modules[peer_spec.name] = peer +peer_spec.loader.exec_module(peer) + class NodeProductCacheTests(unittest.TestCase): def setUp(self): @@ -475,5 +481,503 @@ def test_hit_and_verified_restore_update_bounded_measurement(self): self.assertEqual(hit["snapshot"]["local_hit_rate"], 0.5) + def test_peer_transfer_uses_one_absolute_deadline_across_reads(self): + destination = Path(self.temp.name) / "peer-deadline.tar.gz" + clock = {"now": 0.0} + timeouts = [] + + class FakeSocket: + def settimeout(self, value): + timeouts.append(value) + + class FakeResponse: + status = 200 + + def __init__(self): + self.chunks = [b"a", b"b"] + + def getheader(self, name, default=None): + if name == "Content-Length": + return "2" + return default + + def read(self, _size=-1): + clock["now"] += 0.6 + if self.chunks: + return self.chunks.pop(0) + return b"" + + read1 = read + + class FakeConnection: + def __init__(self): + self.timeout = 1.0 + self.sock = FakeSocket() + self.response = FakeResponse() + + def request(self, *_args, **_kwargs): + return None + + def getresponse(self): + return self.response + + def close(self): + return None + + connection = FakeConnection() + with mock.patch.object( + peer.time, + "monotonic", + side_effect=lambda: clock["now"], + ): + with mock.patch.object( + peer, + "_connection", + return_value=(connection, "peer.example"), + ): + with self.assertRaisesRegex( + peer.PeerUnavailable, + "deadline exceeded", + ): + peer.transfer_http( + peer.PeerSource("https://peer.example"), + "a" * 64, + "read-token", + destination, + 2, + timeout=1.0, + ) + + self.assertEqual(destination.read_bytes(), b"ab") + self.assertGreaterEqual(len(timeouts), 2) + self.assertGreater(timeouts[0], timeouts[-1]) + + def test_peer_transfer_rearms_deadline_before_each_underlying_receive(self): + destination = Path(self.temp.name) / "peer-receive-deadline.tar.gz" + clock = {"now": 0.0} + timeouts = [] + + class FakeSocket: + def settimeout(self, value): + timeouts.append(value) + + class FakeResponse: + status = 200 + + def __init__(self): + self.chunks = [b"a", b"b"] + + def getheader(self, name, default=None): + if name == "Content-Length": + return "2" + return default + + def read(self, _size=-1): + raise AssertionError("buffered response.read must not own the peer deadline") + + def read1(self, _size=-1): + clock["now"] += 0.6 + if self.chunks: + return self.chunks.pop(0) + return b"" + + class FakeConnection: + def __init__(self): + self.timeout = 1.0 + self.sock = FakeSocket() + self.response = FakeResponse() + + def request(self, *_args, **_kwargs): + return None + + def getresponse(self): + return self.response + + def close(self): + return None + + connection = FakeConnection() + with mock.patch.object( + peer.time, + "monotonic", + side_effect=lambda: clock["now"], + ): + with mock.patch.object( + peer, + "_connection", + return_value=(connection, "peer.example"), + ): + with self.assertRaisesRegex( + peer.PeerUnavailable, + "deadline exceeded", + ): + peer.transfer_http( + peer.PeerSource("https://peer.example"), + "a" * 64, + "read-token", + destination, + 2, + timeout=1.0, + ) + + self.assertEqual(destination.read_bytes(), b"ab") + self.assertGreaterEqual(len(timeouts), 2) + self.assertGreater(timeouts[0], timeouts[-1]) + + + def test_peer_server_bounds_client_time_and_active_requests(self): + server = peer.PeerHTTPServer( + ("127.0.0.1", 0), + peer.PeerRequestHandler, + store=self.store, + token="read-token", + drain_marker=None, + client_timeout_seconds=0.25, + max_active_requests=1, + ) + self.addCleanup(server.server_close) + + client = peer.socket.create_connection(server.server_address, timeout=1) + accepted = None + try: + accepted, _ = server.get_request() + self.assertAlmostEqual(accepted.gettimeout(), 0.25) + finally: + if accepted is not None: + accepted.close() + client.close() + + self.assertTrue(server._request_slots.acquire(blocking=False)) + try: + request = mock.Mock() + with mock.patch.object(server, "shutdown_request") as shutdown: + server.process_request(request, ("127.0.0.1", 1)) + shutdown.assert_called_once_with(request) + finally: + server._request_slots.release() + + def test_peer_probe_and_fetch_use_only_exact_object_identity(self): + self.publish() + offer = peer.local_availability(self.store, self.identity.key()) + self.assertIsNotNone(offer) + self.assertEqual(offer.object_key, self.identity.key()) + self.assertEqual(offer.content_digest, self.identity.archive_digest) + self.assertEqual(offer.size_bytes, self.archive.stat().st_size) + + calls = [] + destination = Path(self.temp.name) / "peer-products" + sources = [peer.PeerSource("https://peer-a.example")] + + def probe(source, object_key, token): + calls.append(("probe", source.url, object_key, token)) + return offer + + def transfer(source, object_key, token, target, size): + calls.append(("fetch", source.url, object_key, token, size)) + target.write_bytes(self.archive.read_bytes()) + + result = peer.fetch_exact( + self.identity, + destination, + sources, + probe=probe, + transfer=transfer, + token_loader=lambda _: "read-token", + ) + self.assertTrue(result["hit"], result) + self.assertEqual(result["source"], "peer") + self.assertEqual(result["bytes_transferred"], self.archive.stat().st_size) + self.assertEqual( + (destination / cache.ARCHIVE_NAME).read_bytes(), + self.archive.read_bytes(), + ) + self.assertEqual([call[0] for call in calls], ["probe", "fetch"]) + self.assertTrue(all(call[2] == self.identity.key() for call in calls)) + + def test_corrupt_peer_object_becomes_a_miss_without_partial_destination(self): + self.publish() + offer = peer.local_availability(self.store, self.identity.key()) + destination = Path(self.temp.name) / "corrupt-peer-products" + + def transfer(_source, _key, _token, target, _size): + target.write_bytes(b"corrupt") + + result = peer.fetch_exact( + self.identity, + destination, + [peer.PeerSource("https://peer-a.example")], + probe=lambda *_: offer, + transfer=transfer, + token_loader=lambda _: "read-token", + ) + self.assertFalse(result["hit"], result) + self.assertEqual(result["status"], "miss") + self.assertFalse(destination.exists()) + + def test_peer_death_mid_transfer_falls_through_to_later_peer(self): + self.publish() + offer = peer.local_availability(self.store, self.identity.key()) + destination = Path(self.temp.name) / "peer-failover" + transfers = [] + + def transfer(source, _key, _token, target, _size): + transfers.append(source.url) + if source.url.endswith("peer-a.example"): + target.write_bytes(self.archive.read_bytes()[:64]) + raise TimeoutError("peer disappeared") + target.write_bytes(self.archive.read_bytes()) + + result = peer.fetch_exact( + self.identity, + destination, + [ + peer.PeerSource("https://peer-a.example"), + peer.PeerSource("https://peer-b.example"), + ], + probe=lambda *_: offer, + transfer=transfer, + token_loader=lambda _: "read-token", + ) + self.assertTrue(result["hit"], result) + self.assertEqual(result["source_index"], 1) + self.assertEqual( + transfers, + ["https://peer-a.example", "https://peer-b.example"], + ) + + def test_peer_identity_mismatch_never_uses_filename_equivalence(self): + self.publish() + offer = peer.local_availability(self.store, self.identity.key()) + wrong = peer.PeerAvailability( + object_key="f" * 64, + schema_generation=offer.schema_generation, + size_bytes=offer.size_bytes, + content_digest=offer.content_digest, + ) + destination = Path(self.temp.name) / "wrong-peer-products" + result = peer.fetch_exact( + self.identity, + destination, + [peer.PeerSource("https://peer-a.example")], + probe=lambda *_: wrong, + transfer=lambda *_: self.fail("identity mismatch must refuse before fetch"), + token_loader=lambda _: "read-token", + ) + self.assertFalse(result["hit"], result) + self.assertFalse(destination.exists()) + + def test_consumer_cancellation_cleans_partial_peer_transfer(self): + self.publish() + offer = peer.local_availability(self.store, self.identity.key()) + destination = Path(self.temp.name) / "cancelled-peer-products" + + def transfer(_source, _key, _token, target, _size): + target.write_bytes(self.archive.read_bytes()[:64]) + raise KeyboardInterrupt() + + with self.assertRaises(KeyboardInterrupt): + peer.fetch_exact( + self.identity, + destination, + [peer.PeerSource("https://peer-a.example")], + probe=lambda *_: offer, + transfer=transfer, + token_loader=lambda _: "read-token", + ) + self.assertFalse(destination.exists()) + + def test_source_drain_blocks_new_transfer_but_keeps_existing_lease_valid(self): + self.publish() + draining = {"value": False} + with peer.open_local_object( + self.store, + self.identity.key(), + draining=lambda: draining["value"], + ) as opened: + self.assertEqual(opened.path.read_bytes(), self.archive.read_bytes()) + draining["value"] = True + self.assertTrue(opened.path.exists()) + with self.assertRaises(peer.PeerUnavailable): + with peer.open_local_object( + self.store, + self.identity.key(), + draining=lambda: draining["value"], + ): + pass + result = cache.reclaim(self.store, 0) + self.assertEqual(result["evicted_objects"], 1) + + def test_six_waiters_share_one_peer_fill_and_publish_one_object(self): + owner_destination = Path(self.temp.name) / "peer-owner" + owner = cache.acquire(self.store, self.identity, owner_destination, wait=1) + self.assertTrue(owner["fill"]) + waiter_destinations = [ + Path(self.temp.name) / f"peer-waiter-{index}" for index in range(6) + ] + waiter_results = [None] * 6 + + def waiter(index): + waiter_results[index] = cache.acquire( + self.store, self.identity, waiter_destinations[index], wait=3 + ) + + registrations = 0 + registration_lock = threading.Lock() + all_registered = threading.Event() + register_waiter = cache._register_waiter_locked + + def tracked_register(*args, **kwargs): + nonlocal registrations + result = register_waiter(*args, **kwargs) + if result is not None: + with registration_lock: + registrations += 1 + if registrations == len(waiter_destinations): + all_registered.set() + return result + + threads = [threading.Thread(target=waiter, args=(index,)) for index in range(6)] + with mock.patch.object( + cache, + "_register_waiter_locked", + side_effect=tracked_register, + ): + for thread in threads: + thread.start() + self.assertTrue( + all_registered.wait(3), + "all waiter registrations must exist before peer publication", + ) + + peer_store = cache.Store(Path(self.temp.name) / "peer-source-store") + token = cache.acquire( + peer_store, + self.identity, + Path(self.temp.name) / "peer-source-owner", + wait=0, + )["token"] + seeded = cache.finalize( + peer_store, + self.identity, + self.archive, + token=token, + source_class="github", + restore_succeeded=True, + provider_metadata=lambda _: self.provider(), + ) + self.assertEqual(seeded["status"], "published") + offer = peer.local_availability(peer_store, self.identity.key()) + fetch_calls = [] + + def transfer(_source, _key, _token, target, _size): + fetch_calls.append(1) + with peer.open_local_object( + peer_store, self.identity.key(), draining=lambda: False + ) as opened: + target.write_bytes(opened.path.read_bytes()) + + peer_result = peer.fetch_exact( + self.identity, + owner_destination, + [peer.PeerSource("https://peer-a.example")], + probe=lambda *_: offer, + transfer=transfer, + token_loader=lambda _: "read-token", + ) + self.assertTrue(peer_result["hit"]) + with mock.patch.dict( + os.environ, + { + "GITHUB_REPOSITORY": self.identity.repository, + "GITHUB_RUN_ID": str(self.identity.producer_run_id), + "GITHUB_RUN_ATTEMPT": str(self.identity.producer_run_attempt), + }, + clear=False, + ): + published = cache.finalize( + self.store, + self.identity, + owner_destination / cache.ARCHIVE_NAME, + token=owner["token"], + source_class="peer", + restore_succeeded=True, + provider_metadata=cache.same_run_provider_metadata, + ) + self.assertEqual(published["status"], "published", published) + + for thread in threads: + thread.join(3) + self.assertFalse(thread.is_alive()) + self.assertEqual(len(fetch_calls), 1) + self.assertTrue(all(result["hit"] for result in waiter_results)) + self.assertEqual( + len(list((self.root / "objects" / self.identity.key()[:2]).iterdir())), + 1, + ) + + def test_two_nodes_can_miss_and_install_the_same_peer_object_independently(self): + source_store = cache.Store(Path(self.temp.name) / "peer-source-two-node") + token = cache.acquire( + source_store, + self.identity, + Path(self.temp.name) / "source-fill", + wait=0, + )["token"] + cache.finalize( + source_store, + self.identity, + self.archive, + token=token, + source_class="github", + restore_succeeded=True, + provider_metadata=lambda _: self.provider(), + ) + offer = peer.local_availability(source_store, self.identity.key()) + roots = [ + cache.Store(Path(self.temp.name) / "node-a"), + cache.Store(Path(self.temp.name) / "node-b"), + ] + results = [] + for index, store in enumerate(roots): + destination = Path(self.temp.name) / f"node-{index}-product" + owner = cache.acquire(store, self.identity, destination, wait=0) + self.assertTrue(owner["fill"]) + fetched = peer.fetch_exact( + self.identity, + destination, + [peer.PeerSource("https://peer-a.example")], + probe=lambda *_: offer, + transfer=lambda _source, _key, _token, target, _size: target.write_bytes( + self.archive.read_bytes() + ), + token_loader=lambda _: "read-token", + ) + self.assertTrue(fetched["hit"]) + with mock.patch.dict( + os.environ, + { + "GITHUB_REPOSITORY": self.identity.repository, + "GITHUB_RUN_ID": str(self.identity.producer_run_id), + "GITHUB_RUN_ATTEMPT": str(self.identity.producer_run_attempt), + }, + clear=False, + ): + results.append( + cache.finalize( + store, + self.identity, + destination / cache.ARCHIVE_NAME, + token=owner["token"], + source_class="peer", + restore_succeeded=True, + provider_metadata=cache.same_run_provider_metadata, + ) + ) + self.assertTrue(all(result["status"] == "published" for result in results)) + self.assertTrue( + all(store.entry(self.identity.key()).exists() for store in roots) + ) + + if __name__ == "__main__": unittest.main() diff --git a/workers/ci-artifacts/test/consumer-workflow.test.mjs b/workers/ci-artifacts/test/consumer-workflow.test.mjs index 78a5160b96f1..d900cc3a1af0 100644 --- a/workers/ci-artifacts/test/consumer-workflow.test.mjs +++ b/workers/ci-artifacts/test/consumer-workflow.test.mjs @@ -53,7 +53,6 @@ else: const restore = job.steps.find((step) => step.id === "r2-products"); const fallback = job.steps.find((step) => step.name === "Download compiled app-host test product"); assert.equal(local.run, 'python3 scripts/ci/node_product_cache.py acquire "$RUNNER_TEMP/app-host-products"'); - assert.equal(evaluate(restore.if, { steps: { "node-products": { outputs: { hit: "false" } } } }), "true"); const values = { github: { token: "read-only-job-token" }, vars: { CI_ARTIFACT_R2_URL: enabled ? "https://broker.example" : "" }, @@ -62,13 +61,26 @@ else: artifact_digest: "sha256:169fa1cbfb4a778073f1efdc24904064d3bee6103f4a89e5054936d6661eed2a", } } }, }; + const restoreEnabled = evaluate(restore.if, { + steps: { + "node-products": { outputs: { hit: "false" } }, + "peer-products": { outputs: { hit: "false" } }, + }, + vars: values.vars, + }) === "true"; + assert.equal(restoreEnabled, enabled); const output = path.join(temporary, "output"); - execFileSync("bash", ["-e", "-c", restore.run], { cwd: root, env: { - ...process.env, PATH: `${bin}:${process.env.PATH}`, RUNNER_TEMP: temporary, - GITHUB_OUTPUT: output, GITHUB_RUN_ID: "456", GITHUB_REPOSITORY: "manaflow-ai/cmux", - ...Object.fromEntries(Object.entries(restore.env).map(([key, value]) => [key, render(value, values)])), - } }); - const hit = Object.fromEntries(fs.readFileSync(output, "utf8").trim().split("\n").map((line) => line.split("="))).hit; + let hit = ""; + if (restoreEnabled) { + execFileSync("bash", ["-e", "-c", restore.run], { cwd: root, env: { + ...process.env, PATH: `${bin}:${process.env.PATH}`, RUNNER_TEMP: temporary, + GITHUB_OUTPUT: output, GITHUB_RUN_ID: "456", GITHUB_REPOSITORY: "manaflow-ai/cmux", + ...Object.fromEntries(Object.entries(restore.env).map(([key, value]) => [key, render(value, values)])), + } }); + hit = Object.fromEntries( + fs.readFileSync(output, "utf8").trim().split("\n").map((line) => line.split("=")), + ).hit; + } const download = evaluate(fallback.if, { steps: { "node-products": { outputs: { hit: "false" } }, "r2-products": { outputs: { hit } },