Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
38 changes: 30 additions & 8 deletions amplifier_foundation/modules/activator.py
Original file line number Diff line number Diff line change
Expand Up @@ -574,6 +574,23 @@ async def _install_dependencies(
module_name: str | None = None,
progress_callback: Callable[[str, str], None] | None = None,
force: bool = False,
) -> None:
from .preparation import PreparationReceipt

receipt = PreparationReceipt(self, module_path)
if not force and await asyncio.to_thread(receipt.reusable):
return
await self._install_dependencies_uncached(
module_path, module_name, progress_callback, force
)
await asyncio.to_thread(receipt.record)

async def _install_dependencies_uncached(
self,
module_path: Path,
module_name: str | None = None,
progress_callback: Callable[[str, str], None] | None = None,
force: bool = False,
) -> None:
"""Install Python dependencies for a module.

Expand Down Expand Up @@ -675,6 +692,13 @@ async def _install_dependencies(
requirements = module_path / "requirements.txt"

if pyproject.exists():
from amplifier_foundation.sources.shared import build_lock, build_view

install_path, immutable = await asyncio.to_thread(build_view, module_path)
# Shared source trees are read-only. Build wheels from a private
# writable input; never create an editable installation on them.
# uv reuses the wheel for the same source and build contract.
# Build metadata may change without rewriting the canonical source.
# Build overrides for git URL dependencies that are already installed.
# This prevents uv from fetching/building packages from git when a
# prebuilt wheel is already available (e.g. amplifier-core from PyPI).
Expand All @@ -689,8 +713,8 @@ async def _install_dependencies(
"uv",
"pip",
"install",
"-e",
str(module_path),
*([] if immutable else ["-e"]),
str(install_path),
"--python",
self.install_python,
"--quiet",
Expand Down Expand Up @@ -724,12 +748,10 @@ async def _install_dependencies(
overrides_file.close()
cmd.extend(["--overrides", overrides_file.name])

subprocess.run(
cmd,
check=True,
capture_output=True,
text=True,
)
# Holding this across uv prevents concurrent build backends
# from modifying the same writable source view.
with build_lock(install_path):
subprocess.run(cmd, check=True, capture_output=True, text=True)
# Mark as installed after successful install
self._install_state.mark_installed(module_path)
# Refresh Python's package discovery so subprocess-installed packages
Expand Down
104 changes: 104 additions & 0 deletions amplifier_foundation/modules/preparation.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,104 @@
"""Reuse a completed install only inside one isolated preparation attempt.

Every new attempt refreshes each selected source at least once. A duplicate
may skip resolution only when the source, explicit policy and complete installed
package metadata still match. Any intervening graph change forces a fresh install.
This is not a standing installed-version shortcut or a dependency pin.
"""

from __future__ import annotations

import hashlib
import importlib.metadata
import json
import os
import re
import sys
from pathlib import Path


def graph_signature():
rows = []
for dist in importlib.metadata.distributions():
rows.append(
[
dist.metadata["Name"],
dist.version,
dist.read_text("METADATA"),
dist.read_text("direct_url.json"),
]
)
return hashlib.sha256(json.dumps(sorted(rows), sort_keys=True).encode()).hexdigest()


def source_signature(root):
root = Path(root)
digest = hashlib.sha256()
for path in sorted(root.rglob("*")):
relative = path.relative_to(root)
if any(
part in {".git", ".venv", "node_modules", "__pycache__", "build", "dist"}
or part.endswith(".egg-info")
for part in relative.parts
):
continue
if path.is_symlink():
if not path.resolve().is_relative_to(root.resolve()):
# External build inputs can change independently. Unknown
# content must force installation, not a reusable receipt.
raise ValueError("External source symlink cannot be qualified")
digest.update(str(relative).encode() + b"\0" + os.readlink(path).encode())
elif path.is_file():
digest.update(str(relative).encode() + b"\0" + path.read_bytes())
return digest.hexdigest()


class PreparationReceipt:
def __init__(self, activator, source):
self.source = Path(source)
token = os.environ.get("AMPLIFIER_INSTALL_PREPARATION", "")
# An external install target cannot be qualified by this process's graph.
self.path = None
if (
activator.refresh_dependencies
and re.fullmatch("[a-f0-9]{32}", token)
and os.path.abspath(activator.install_python)
== os.path.abspath(sys.executable)
):
policy = [
Path(value).read_bytes().hex() if value else None
for value in (
activator.install_constraints,
activator.install_overrides,
)
]
self.policy = hashlib.sha256(json.dumps(policy).encode()).hexdigest()
key = hashlib.sha256(str(self.source.absolute()).encode()).hexdigest()
self.path = (
activator.cache_dir / "install-preparations" / token / (key + ".json")
)

def evidence(self):
return {
"source": source_signature(self.source),
"policy": self.policy,
"graph": graph_signature(),
}

def reusable(self):
if not self.path:
return False
try:
return json.loads(self.path.read_text()) == self.evidence()
except (OSError, ValueError, TypeError):
return False

def record(self):
if self.path:
from amplifier_foundation.settings import atomic_write

try:
self.path.parent.mkdir(parents=True, exist_ok=True)
atomic_write(self.path, json.dumps(self.evidence()), private=True)
except (OSError, ValueError, TypeError):
pass # Reuse is optional; later qualification remains authoritative.
10 changes: 8 additions & 2 deletions amplifier_foundation/registry.py
Original file line number Diff line number Diff line change
Expand Up @@ -185,6 +185,7 @@ def __init__(
strict: bool = False,
include_source_resolver: Callable[[str], str | None] | None = None,
persist: bool = True,
read_persisted: bool = True,
) -> None:
"""Initialize registry.

Expand All @@ -204,6 +205,10 @@ def __init__(
cleanup and explicit save() calls cannot write registry.json.
Source downloads still use the shared content cache. Use a
fresh instance per session with scoped source overrides.
read_persisted: If False, start with no saved registrations. The
caller supplies its authoritative registrations; the source
cache still belongs to the explicit home. Defaults to True
so existing CLI and library consumers retain their behavior.
"""
self._home = self._resolve_home(home)
self._strict = strict
Expand All @@ -217,8 +222,9 @@ def __init__(
# Future-based deduplication: cache loaded bundles and track in-progress loads
self._loaded_bundles: dict[str, Bundle] = {} # Cache of fully loaded bundles
self._pending_loads: dict[str, asyncio.Future[Bundle]] = {} # In-progress loads
self._load_persisted_state()
self._validate_cached_paths()
if read_persisted:
self._load_persisted_state()
self._validate_cached_paths()

@property
def home(self) -> Path:
Expand Down
60 changes: 59 additions & 1 deletion amplifier_foundation/sources/git.py
Original file line number Diff line number Diff line change
Expand Up @@ -411,6 +411,14 @@ def _build_git_url(self, parsed: ParsedURI) -> str:
scheme = parsed.scheme.replace("git+", "")
return f"{scheme}://{parsed.host}{parsed.path}"

def _shared_uri(self, parsed):
uri = "git+" + self._build_git_url(parsed)
if parsed.ref:
uri += "@" + parsed.ref
if parsed.subpath:
uri += "#subdirectory=" + parsed.subpath
return uri

def _get_cache_path(self, parsed: ParsedURI, cache_dir: Path) -> Path:
"""Get the cache path for a parsed URI."""
git_url = self._build_git_url(parsed)
Expand Down Expand Up @@ -672,6 +680,20 @@ async def resolve(self, parsed: ParsedURI, cache_dir: Path) -> ResolvedSource:
Raises:
BundleNotFoundError: If clone fails or ref not found.
"""
if os.environ.get("AMPLIFIER_SOURCE_STORE"):
from .shared import resolve_shared_source

result = await resolve_shared_source(
parsed.original
if hasattr(parsed, "original")
else self._shared_uri(parsed),
cache_dir,
)
if not self._verify_clone_integrity(result.source_root):
raise BundleNotFoundError(
"Shared source is not a bundle or module repository"
)
return result
cache_path = self._get_cache_path(parsed, cache_dir)
cache_path.parent.mkdir(parents=True, exist_ok=True)
# Workers share this cache across processes. Check integrity only after
Expand Down Expand Up @@ -829,6 +851,31 @@ async def get_status(self, parsed: ParsedURI, cache_dir: Path) -> SourceStatus:
source_uri += f"#subdirectory={parsed.subpath}"

# Initialize status
if shared_root := os.environ.get("AMPLIFIER_SOURCE_STORE"):
from .shared import SharedSourceStore

try:
shared_path = await asyncio.to_thread(
SharedSourceStore(shared_root).cached, source_uri, cache_dir
)
if shared_path is not None:
cache_path = shared_path
except (
OSError,
ValueError,
KeyError,
TypeError,
subprocess.SubprocessError,
):
return SourceStatus(
source_uri=source_uri,
is_cached=False,
cached_ref=ref,
remote_ref=ref,
has_update=None,
error="Shared source could not be verified",
summary="Shared source was retained for inspection",
)
status = SourceStatus(
source_uri=source_uri,
is_cached=cache_path.exists(),
Expand Down Expand Up @@ -856,7 +903,12 @@ async def get_status(self, parsed: ParsedURI, cache_dir: Path) -> SourceStatus:

# Get remote commit
try:
status.remote_commit = await self._get_remote_commit(git_url, ref)
if shared_root:
from .shared import remote_revision

status.remote_commit = await remote_revision(git_url, ref)
else:
status.remote_commit = await self._get_remote_commit(git_url, ref)

if status.remote_commit is None:
status.has_update = None
Expand Down Expand Up @@ -900,6 +952,12 @@ async def update(self, parsed: ParsedURI, cache_dir: Path) -> ResolvedSource:
Raises:
BundleNotFoundError: If clone fails.
"""
if os.environ.get("AMPLIFIER_SOURCE_STORE"):
from .shared import resolve_shared_source

return await resolve_shared_source(
self._shared_uri(parsed), cache_dir, refresh=True
)
cache_path = self._get_cache_path(parsed, cache_dir)
cache_path.parent.mkdir(parents=True, exist_ok=True)
async with AsyncFileLock(cache_path.with_name(f".{cache_path.name}.lock")):
Expand Down
Loading
Loading