977 lines
38 KiB
Python
977 lines
38 KiB
Python
"""Optional content-addressed artifact exchange with a bounded local cache."""
|
|
|
|
from __future__ import annotations
|
|
|
|
import fcntl
|
|
import hashlib
|
|
import json
|
|
import os
|
|
import re
|
|
import shutil
|
|
import sqlite3
|
|
import stat
|
|
import threading
|
|
import time
|
|
from collections.abc import Iterator, Mapping, Sequence
|
|
from contextlib import contextmanager, suppress
|
|
from dataclasses import dataclass
|
|
from pathlib import Path
|
|
from typing import Any, Final, Literal
|
|
from uuid import uuid4
|
|
|
|
from k1link.artifacts import utc_now_iso
|
|
|
|
MANIFEST_SCHEMA: Final = "missioncore.artifact-manifest/v1"
|
|
REFERENCE_SCHEMA: Final = "missioncore.artifact-reference/v1"
|
|
DEFAULT_CACHE_MAX_BYTES: Final = 8 * 1024 * 1024 * 1024
|
|
DEFAULT_FREE_SPACE_RESERVE_BYTES: Final = 2 * 1024 * 1024 * 1024
|
|
STORE_ROOT_ENV: Final = "MISSIONCORE_ARTIFACT_STORE_ROOT"
|
|
CACHE_MAX_BYTES_ENV: Final = "MISSIONCORE_ARTIFACT_CACHE_MAX_BYTES"
|
|
CACHE_RESERVE_BYTES_ENV: Final = "MISSIONCORE_ARTIFACT_CACHE_FREE_SPACE_RESERVE_BYTES"
|
|
|
|
_SHA256 = re.compile(r"^[a-f0-9]{64}$")
|
|
_SAFE_COMPONENT = re.compile(r"^[A-Za-z0-9][A-Za-z0-9._:-]{0,127}$")
|
|
_MAX_MANIFEST_BYTES = 8 * 1024 * 1024
|
|
_COPY_CHUNK_BYTES = 1024 * 1024
|
|
|
|
|
|
class ArtifactGatewayError(RuntimeError):
|
|
"""Base class for central artifact and host-cache failures."""
|
|
|
|
|
|
class ArtifactStoreUnavailable(ArtifactGatewayError):
|
|
"""The configured central store cannot currently be reached."""
|
|
|
|
|
|
class ArtifactNotFound(ArtifactGatewayError):
|
|
"""An exact reference, manifest, role, or object does not exist."""
|
|
|
|
|
|
class ArtifactIntegrityError(ArtifactGatewayError):
|
|
"""An artifact violates its content-addressed identity."""
|
|
|
|
|
|
class ArtifactCacheCapacityError(ArtifactGatewayError):
|
|
"""The local host cannot safely reserve enough cache capacity."""
|
|
|
|
|
|
@dataclass(frozen=True, slots=True)
|
|
class ArtifactMember:
|
|
role: str
|
|
media_type: str
|
|
sha256: str
|
|
byte_length: int
|
|
|
|
|
|
@dataclass(frozen=True, slots=True)
|
|
class ArtifactManifest:
|
|
manifest_id: str
|
|
artifact_type: str
|
|
subject_id: str
|
|
created_at_utc: str
|
|
members: tuple[ArtifactMember, ...]
|
|
metadata: Mapping[str, str]
|
|
|
|
def member(self, role: str) -> ArtifactMember:
|
|
matches = tuple(item for item in self.members if item.role == role)
|
|
if len(matches) != 1:
|
|
raise ArtifactNotFound(f"artifact role is unavailable: {role}")
|
|
return matches[0]
|
|
|
|
|
|
@dataclass(frozen=True, slots=True)
|
|
class ResolvedArtifact:
|
|
manifest: ArtifactManifest
|
|
member: ArtifactMember
|
|
path: Path
|
|
cache_hit: bool
|
|
central_available: bool
|
|
|
|
|
|
@dataclass(frozen=True, slots=True)
|
|
class ArtifactCacheStatus:
|
|
object_count: int
|
|
total_bytes: int
|
|
pinned_object_count: int
|
|
pinned_bytes: int
|
|
cache_max_bytes: int
|
|
free_space_reserve_bytes: int
|
|
|
|
|
|
@dataclass(frozen=True, slots=True)
|
|
class ArtifactGatewayStatus:
|
|
central_status: Literal["ready", "unavailable", "invalid"]
|
|
cache: ArtifactCacheStatus
|
|
|
|
|
|
class CentralArtifactStore:
|
|
"""Immutable SHA-256 objects/manifests plus atomically replaceable named refs."""
|
|
|
|
def __init__(self, root: Path, *, create: bool = False) -> None:
|
|
self.root = root.expanduser().absolute()
|
|
if create:
|
|
try:
|
|
self.root.mkdir(parents=True, exist_ok=True)
|
|
except OSError as exc:
|
|
raise ArtifactStoreUnavailable("central artifact store cannot be created") from exc
|
|
self._require_root()
|
|
|
|
def publish_file(self, source: Path) -> ArtifactMember:
|
|
source_path = _regular_source(source)
|
|
digest, byte_length = _hash_file(source_path)
|
|
destination = self.object_path(digest)
|
|
try:
|
|
destination.parent.mkdir(parents=True, exist_ok=True)
|
|
if destination.exists() or destination.is_symlink():
|
|
_verify_object(destination, digest, byte_length)
|
|
else:
|
|
temporary = destination.with_name(f".{digest}.{uuid4().hex}.tmp")
|
|
try:
|
|
copied_digest, copied_bytes = _copy_and_hash(source_path, temporary)
|
|
if copied_digest != digest or copied_bytes != byte_length:
|
|
raise ArtifactIntegrityError("source changed while it was published")
|
|
os.replace(temporary, destination)
|
|
finally:
|
|
temporary.unlink(missing_ok=True)
|
|
except ArtifactGatewayError:
|
|
raise
|
|
except OSError as exc:
|
|
raise ArtifactStoreUnavailable("central artifact object cannot be published") from exc
|
|
return ArtifactMember(
|
|
role="",
|
|
media_type="application/octet-stream",
|
|
sha256=digest,
|
|
byte_length=byte_length,
|
|
)
|
|
|
|
def publish_manifest(
|
|
self,
|
|
*,
|
|
artifact_type: str,
|
|
subject_id: str,
|
|
members: Sequence[ArtifactMember],
|
|
metadata: Mapping[str, str] | None = None,
|
|
created_at_utc: str | None = None,
|
|
) -> ArtifactManifest:
|
|
_validate_component(artifact_type, "artifact type")
|
|
_validate_component(subject_id, "artifact subject")
|
|
normalized_members = tuple(sorted(members, key=lambda item: item.role))
|
|
_validate_members(normalized_members)
|
|
normalized_metadata = dict(sorted((metadata or {}).items()))
|
|
for key, value in normalized_metadata.items():
|
|
_validate_component(key, "metadata key")
|
|
if not isinstance(value, str) or len(value) > 4096:
|
|
raise ValueError("artifact metadata value is invalid")
|
|
document = {
|
|
"schema_version": MANIFEST_SCHEMA,
|
|
"artifact_type": artifact_type,
|
|
"subject_id": subject_id,
|
|
"created_at_utc": created_at_utc or utc_now_iso(),
|
|
"members": [
|
|
{
|
|
"role": item.role,
|
|
"media_type": item.media_type,
|
|
"sha256": item.sha256,
|
|
"byte_length": item.byte_length,
|
|
}
|
|
for item in normalized_members
|
|
],
|
|
"metadata": normalized_metadata,
|
|
}
|
|
payload = _canonical_json(document)
|
|
manifest_id = hashlib.sha256(payload).hexdigest()
|
|
destination = self.manifest_path(manifest_id)
|
|
try:
|
|
destination.parent.mkdir(parents=True, exist_ok=True)
|
|
if destination.exists() or destination.is_symlink():
|
|
if destination.read_bytes() != payload:
|
|
raise ArtifactIntegrityError("central artifact manifest digest is inconsistent")
|
|
else:
|
|
_write_bytes_atomic(destination, payload)
|
|
except ArtifactGatewayError:
|
|
raise
|
|
except OSError as exc:
|
|
raise ArtifactStoreUnavailable("central artifact manifest cannot be published") from exc
|
|
return _parse_manifest(document, manifest_id)
|
|
|
|
def set_reference(self, namespace: str, key: str, manifest_id: str) -> None:
|
|
_validate_component(namespace, "artifact namespace")
|
|
_validate_component(key, "artifact reference key")
|
|
_validate_sha256(manifest_id)
|
|
self.read_manifest(manifest_id)
|
|
document = {
|
|
"schema_version": REFERENCE_SCHEMA,
|
|
"namespace": namespace,
|
|
"key": key,
|
|
"manifest_id": manifest_id,
|
|
"updated_at_utc": utc_now_iso(),
|
|
}
|
|
try:
|
|
_write_bytes_atomic(self.reference_path(namespace, key), _canonical_json(document))
|
|
except OSError as exc:
|
|
raise ArtifactStoreUnavailable(
|
|
"central artifact reference cannot be published"
|
|
) from exc
|
|
|
|
def resolve_reference(self, namespace: str, key: str) -> ArtifactManifest:
|
|
reference = self.read_reference_document(namespace, key)
|
|
return self.read_manifest(str(reference["manifest_id"]))
|
|
|
|
def read_reference_document(self, namespace: str, key: str) -> dict[str, Any]:
|
|
_validate_component(namespace, "artifact namespace")
|
|
_validate_component(key, "artifact reference key")
|
|
path = self.reference_path(namespace, key)
|
|
document = _read_json_document(path, unavailable_message="central reference is unavailable")
|
|
if (
|
|
document.get("schema_version") != REFERENCE_SCHEMA
|
|
or document.get("namespace") != namespace
|
|
or document.get("key") != key
|
|
or not isinstance(document.get("manifest_id"), str)
|
|
or _SHA256.fullmatch(str(document["manifest_id"])) is None
|
|
or not isinstance(document.get("updated_at_utc"), str)
|
|
):
|
|
raise ArtifactIntegrityError("central artifact reference is invalid")
|
|
return document
|
|
|
|
def read_manifest(self, manifest_id: str) -> ArtifactManifest:
|
|
_validate_sha256(manifest_id)
|
|
path = self.manifest_path(manifest_id)
|
|
document = _read_json_document(path, unavailable_message="central manifest is unavailable")
|
|
payload = _canonical_json(document)
|
|
if hashlib.sha256(payload).hexdigest() != manifest_id:
|
|
raise ArtifactIntegrityError("central artifact manifest digest changed")
|
|
return _parse_manifest(document, manifest_id)
|
|
|
|
def object_path(self, digest: str) -> Path:
|
|
_validate_sha256(digest)
|
|
return self.root / "objects" / "sha256" / digest[:2] / digest
|
|
|
|
def manifest_path(self, manifest_id: str) -> Path:
|
|
_validate_sha256(manifest_id)
|
|
return self.root / "manifests" / "sha256" / manifest_id[:2] / f"{manifest_id}.json"
|
|
|
|
def reference_path(self, namespace: str, key: str) -> Path:
|
|
_validate_component(namespace, "artifact namespace")
|
|
_validate_component(key, "artifact reference key")
|
|
return self.root / "refs" / namespace / f"{key}.json"
|
|
|
|
def _require_root(self) -> None:
|
|
try:
|
|
metadata = self.root.lstat()
|
|
except OSError as exc:
|
|
raise ArtifactStoreUnavailable("central artifact store is unavailable") from exc
|
|
if stat.S_ISLNK(metadata.st_mode) or not stat.S_ISDIR(metadata.st_mode):
|
|
raise ArtifactIntegrityError("central artifact store root is invalid")
|
|
|
|
|
|
class LocalArtifactCache:
|
|
"""Persistent local CAS with durable pins and unpinned LRU eviction."""
|
|
|
|
def __init__(
|
|
self,
|
|
root: Path,
|
|
*,
|
|
max_bytes: int = DEFAULT_CACHE_MAX_BYTES,
|
|
free_space_reserve_bytes: int = DEFAULT_FREE_SPACE_RESERVE_BYTES,
|
|
) -> None:
|
|
if max_bytes <= 0:
|
|
raise ValueError("artifact cache maximum bytes must be positive")
|
|
if free_space_reserve_bytes < 0:
|
|
raise ValueError("artifact cache free-space reserve must be non-negative")
|
|
self.root = root.expanduser().absolute()
|
|
self.root.mkdir(mode=0o700, parents=True, exist_ok=True)
|
|
root_metadata = self.root.lstat()
|
|
if stat.S_ISLNK(root_metadata.st_mode) or not stat.S_ISDIR(root_metadata.st_mode):
|
|
raise ArtifactIntegrityError("local artifact cache root is invalid")
|
|
with suppress(OSError):
|
|
self.root.chmod(0o700)
|
|
self.objects_root = self.root / "objects" / "sha256"
|
|
self.metadata_root = self.root / "metadata"
|
|
self.objects_root.mkdir(mode=0o700, parents=True, exist_ok=True)
|
|
self.metadata_root.mkdir(mode=0o700, parents=True, exist_ok=True)
|
|
self.database_path = self.root / "cache.sqlite3"
|
|
self.lock_path = self.root / ".cache.lock"
|
|
self.max_bytes = max_bytes
|
|
self.free_space_reserve_bytes = free_space_reserve_bytes
|
|
self._thread_lock = threading.RLock()
|
|
self._initialize_database()
|
|
|
|
def get(self, member: ArtifactMember) -> Path | None:
|
|
with self._locked():
|
|
return self._get_locked(member, touch=True)
|
|
|
|
def fetch(
|
|
self,
|
|
store: CentralArtifactStore,
|
|
member: ArtifactMember,
|
|
*,
|
|
pin_id: str | None = None,
|
|
) -> tuple[Path, bool]:
|
|
_validate_member(member)
|
|
if pin_id is not None:
|
|
_validate_component(pin_id, "artifact pin id")
|
|
with self._locked():
|
|
cached = self._get_locked(member, touch=True)
|
|
if cached is not None:
|
|
if pin_id is not None:
|
|
self._pin_locked(pin_id, member.sha256)
|
|
return cached, True
|
|
self._reserve_locked(member.byte_length, protected_digest=member.sha256)
|
|
source = store.object_path(member.sha256)
|
|
destination = self.object_path(member.sha256)
|
|
destination.parent.mkdir(mode=0o700, parents=True, exist_ok=True)
|
|
temporary = destination.with_name(f".{member.sha256}.{uuid4().hex}.tmp")
|
|
try:
|
|
copied_digest, copied_bytes = _copy_and_hash(source, temporary)
|
|
if copied_digest != member.sha256 or copied_bytes != member.byte_length:
|
|
raise ArtifactIntegrityError("central artifact object failed verification")
|
|
os.chmod(temporary, 0o600)
|
|
os.replace(temporary, destination)
|
|
except ArtifactGatewayError:
|
|
raise
|
|
except OSError as exc:
|
|
raise ArtifactStoreUnavailable("central artifact object is unavailable") from exc
|
|
finally:
|
|
temporary.unlink(missing_ok=True)
|
|
metadata = destination.stat()
|
|
with self._connect() as connection:
|
|
connection.execute(
|
|
"""
|
|
INSERT INTO objects(
|
|
sha256, byte_length, last_access_ns, mtime_ns, inode
|
|
) VALUES (?, ?, ?, ?, ?)
|
|
ON CONFLICT(sha256) DO UPDATE SET
|
|
byte_length=excluded.byte_length,
|
|
last_access_ns=excluded.last_access_ns,
|
|
mtime_ns=excluded.mtime_ns,
|
|
inode=excluded.inode
|
|
""",
|
|
(
|
|
member.sha256,
|
|
member.byte_length,
|
|
time.time_ns(),
|
|
metadata.st_mtime_ns,
|
|
metadata.st_ino,
|
|
),
|
|
)
|
|
if pin_id is not None:
|
|
connection.execute(
|
|
"INSERT OR IGNORE INTO pins(pin_id, sha256) VALUES (?, ?)",
|
|
(pin_id, member.sha256),
|
|
)
|
|
self._evict_locked(protected_digest=member.sha256)
|
|
return destination.resolve(strict=True), False
|
|
|
|
def pin(self, pin_id: str, members: Sequence[ArtifactMember]) -> tuple[Path, ...]:
|
|
_validate_component(pin_id, "artifact pin id")
|
|
paths: list[Path] = []
|
|
try:
|
|
for member in members:
|
|
path = self.get(member)
|
|
if path is None:
|
|
raise ArtifactNotFound(
|
|
f"artifact object must be fetched before pinning: {member.role}"
|
|
)
|
|
with self._locked():
|
|
self._pin_locked(pin_id, member.sha256)
|
|
paths.append(path)
|
|
except BaseException:
|
|
self.unpin(pin_id)
|
|
raise
|
|
return tuple(paths)
|
|
|
|
def unpin(self, pin_id: str) -> None:
|
|
_validate_component(pin_id, "artifact pin id")
|
|
with self._locked(), self._connect() as connection:
|
|
connection.execute("DELETE FROM pins WHERE pin_id = ?", (pin_id,))
|
|
with self._locked():
|
|
self._evict_locked(protected_digest=None)
|
|
|
|
def save_snapshot(
|
|
self,
|
|
namespace: str,
|
|
key: str,
|
|
reference: Mapping[str, Any],
|
|
manifest: ArtifactManifest,
|
|
) -> None:
|
|
_validate_component(namespace, "artifact namespace")
|
|
_validate_component(key, "artifact reference key")
|
|
manifest_document = _manifest_document(manifest)
|
|
reference_document = dict(reference)
|
|
_validate_reference_snapshot(reference_document, namespace, key)
|
|
if reference_document["manifest_id"] != manifest.manifest_id:
|
|
raise ArtifactIntegrityError("artifact reference and manifest do not match")
|
|
manifest_path = (
|
|
self.metadata_root
|
|
/ "manifests"
|
|
/ manifest.manifest_id[:2]
|
|
/ f"{manifest.manifest_id}.json"
|
|
)
|
|
reference_path = self.metadata_root / "refs" / namespace / f"{key}.json"
|
|
with self._locked():
|
|
_write_bytes_atomic(manifest_path, _canonical_json(manifest_document), mode=0o600)
|
|
_write_bytes_atomic(reference_path, _canonical_json(reference_document), mode=0o600)
|
|
|
|
def load_snapshot(self, namespace: str, key: str) -> ArtifactManifest:
|
|
_validate_component(namespace, "artifact namespace")
|
|
_validate_component(key, "artifact reference key")
|
|
reference_path = self.metadata_root / "refs" / namespace / f"{key}.json"
|
|
try:
|
|
reference = _read_json_document(
|
|
reference_path,
|
|
unavailable_message="local artifact reference is unavailable",
|
|
)
|
|
except ArtifactStoreUnavailable as exc:
|
|
raise ArtifactNotFound("local artifact reference is unavailable") from exc
|
|
_validate_reference_snapshot(reference, namespace, key)
|
|
manifest_id = str(reference["manifest_id"])
|
|
manifest_path = (
|
|
self.metadata_root / "manifests" / manifest_id[:2] / f"{manifest_id}.json"
|
|
)
|
|
try:
|
|
document = _read_json_document(
|
|
manifest_path,
|
|
unavailable_message="local artifact manifest is unavailable",
|
|
)
|
|
except ArtifactStoreUnavailable as exc:
|
|
raise ArtifactNotFound("local artifact manifest is unavailable") from exc
|
|
if hashlib.sha256(_canonical_json(document)).hexdigest() != manifest_id:
|
|
raise ArtifactIntegrityError("local artifact manifest digest changed")
|
|
return _parse_manifest(document, manifest_id)
|
|
|
|
def status(self) -> ArtifactCacheStatus:
|
|
with self._locked(), self._connect() as connection:
|
|
row = connection.execute(
|
|
"""
|
|
SELECT
|
|
COUNT(*),
|
|
COALESCE(SUM(byte_length), 0),
|
|
COUNT(DISTINCT pins.sha256),
|
|
COALESCE(SUM(
|
|
CASE WHEN pins.sha256 IS NOT NULL THEN objects.byte_length ELSE 0 END
|
|
), 0)
|
|
FROM objects
|
|
LEFT JOIN (SELECT DISTINCT sha256 FROM pins) AS pins USING (sha256)
|
|
"""
|
|
).fetchone()
|
|
assert row is not None
|
|
return ArtifactCacheStatus(
|
|
object_count=int(row[0]),
|
|
total_bytes=int(row[1]),
|
|
pinned_object_count=int(row[2]),
|
|
pinned_bytes=int(row[3]),
|
|
cache_max_bytes=self.max_bytes,
|
|
free_space_reserve_bytes=self.free_space_reserve_bytes,
|
|
)
|
|
|
|
def object_path(self, digest: str) -> Path:
|
|
_validate_sha256(digest)
|
|
return self.objects_root / digest[:2] / digest
|
|
|
|
def _initialize_database(self) -> None:
|
|
with self._connect() as connection:
|
|
connection.executescript(
|
|
"""
|
|
PRAGMA journal_mode = WAL;
|
|
PRAGMA foreign_keys = ON;
|
|
CREATE TABLE IF NOT EXISTS objects (
|
|
sha256 TEXT PRIMARY KEY,
|
|
byte_length INTEGER NOT NULL CHECK (byte_length >= 0),
|
|
last_access_ns INTEGER NOT NULL,
|
|
mtime_ns INTEGER NOT NULL,
|
|
inode INTEGER NOT NULL
|
|
);
|
|
CREATE TABLE IF NOT EXISTS pins (
|
|
pin_id TEXT NOT NULL,
|
|
sha256 TEXT NOT NULL REFERENCES objects(sha256) ON DELETE CASCADE,
|
|
PRIMARY KEY(pin_id, sha256)
|
|
);
|
|
CREATE INDEX IF NOT EXISTS objects_lru ON objects(last_access_ns, sha256);
|
|
CREATE INDEX IF NOT EXISTS pins_digest ON pins(sha256);
|
|
"""
|
|
)
|
|
with suppress(OSError):
|
|
self.database_path.chmod(0o600)
|
|
|
|
def _connect(self) -> sqlite3.Connection:
|
|
connection = sqlite3.connect(self.database_path, timeout=30.0)
|
|
connection.execute("PRAGMA foreign_keys = ON")
|
|
return connection
|
|
|
|
@contextmanager
|
|
def _locked(self) -> Iterator[None]:
|
|
with self._thread_lock:
|
|
descriptor = os.open(self.lock_path, os.O_CREAT | os.O_RDWR, 0o600)
|
|
try:
|
|
fcntl.flock(descriptor, fcntl.LOCK_EX)
|
|
yield
|
|
finally:
|
|
fcntl.flock(descriptor, fcntl.LOCK_UN)
|
|
os.close(descriptor)
|
|
|
|
def _get_locked(self, member: ArtifactMember, *, touch: bool) -> Path | None:
|
|
_validate_member(member)
|
|
path = self.object_path(member.sha256)
|
|
with self._connect() as connection:
|
|
row = connection.execute(
|
|
"SELECT byte_length, mtime_ns, inode FROM objects WHERE sha256 = ?",
|
|
(member.sha256,),
|
|
).fetchone()
|
|
if row is None:
|
|
return None
|
|
try:
|
|
metadata = path.lstat()
|
|
except OSError:
|
|
connection.execute("DELETE FROM objects WHERE sha256 = ?", (member.sha256,))
|
|
return None
|
|
if (
|
|
stat.S_ISLNK(metadata.st_mode)
|
|
or not stat.S_ISREG(metadata.st_mode)
|
|
or metadata.st_size != member.byte_length
|
|
or int(row[0]) != member.byte_length
|
|
):
|
|
path.unlink(missing_ok=True)
|
|
connection.execute("DELETE FROM objects WHERE sha256 = ?", (member.sha256,))
|
|
return None
|
|
if metadata.st_mtime_ns != int(row[1]) or metadata.st_ino != int(row[2]):
|
|
digest, byte_length = _hash_file(path)
|
|
if digest != member.sha256 or byte_length != member.byte_length:
|
|
path.unlink(missing_ok=True)
|
|
connection.execute("DELETE FROM objects WHERE sha256 = ?", (member.sha256,))
|
|
return None
|
|
connection.execute(
|
|
"UPDATE objects SET mtime_ns = ?, inode = ? WHERE sha256 = ?",
|
|
(metadata.st_mtime_ns, metadata.st_ino, member.sha256),
|
|
)
|
|
if touch:
|
|
connection.execute(
|
|
"UPDATE objects SET last_access_ns = ? WHERE sha256 = ?",
|
|
(time.time_ns(), member.sha256),
|
|
)
|
|
return path.resolve(strict=True)
|
|
|
|
def _pin_locked(self, pin_id: str, digest: str) -> None:
|
|
with self._connect() as connection:
|
|
connection.execute(
|
|
"INSERT OR IGNORE INTO pins(pin_id, sha256) VALUES (?, ?)",
|
|
(pin_id, digest),
|
|
)
|
|
|
|
def _reserve_locked(self, required_bytes: int, *, protected_digest: str) -> None:
|
|
if required_bytes < 0:
|
|
raise ArtifactCacheCapacityError("artifact cache reservation is invalid")
|
|
if required_bytes > self.max_bytes:
|
|
raise ArtifactCacheCapacityError("artifact exceeds the local cache quota")
|
|
self._evict_locked(
|
|
protected_digest=protected_digest,
|
|
incoming_bytes=required_bytes,
|
|
)
|
|
with self._connect() as connection:
|
|
row = connection.execute(
|
|
"SELECT COALESCE(SUM(byte_length), 0) FROM objects"
|
|
).fetchone()
|
|
total_bytes = int(row[0]) if row is not None else 0
|
|
free_bytes = shutil.disk_usage(self.root).free
|
|
if (
|
|
total_bytes + required_bytes > self.max_bytes
|
|
or free_bytes - required_bytes < self.free_space_reserve_bytes
|
|
):
|
|
raise ArtifactCacheCapacityError(
|
|
"local artifact cache cannot preserve its quota and free-space reserve"
|
|
)
|
|
|
|
def _evict_locked(
|
|
self,
|
|
*,
|
|
protected_digest: str | None,
|
|
incoming_bytes: int = 0,
|
|
) -> None:
|
|
while True:
|
|
with self._connect() as connection:
|
|
total_row = connection.execute(
|
|
"SELECT COALESCE(SUM(byte_length), 0) FROM objects"
|
|
).fetchone()
|
|
total_bytes = int(total_row[0]) if total_row is not None else 0
|
|
free_bytes = shutil.disk_usage(self.root).free
|
|
if (
|
|
total_bytes + incoming_bytes <= self.max_bytes
|
|
and free_bytes - incoming_bytes >= self.free_space_reserve_bytes
|
|
):
|
|
return
|
|
row = connection.execute(
|
|
"""
|
|
SELECT objects.sha256
|
|
FROM objects
|
|
LEFT JOIN pins ON pins.sha256 = objects.sha256
|
|
WHERE pins.sha256 IS NULL AND objects.sha256 != COALESCE(?, '')
|
|
ORDER BY objects.last_access_ns ASC, objects.sha256 ASC
|
|
LIMIT 1
|
|
""",
|
|
(protected_digest,),
|
|
).fetchone()
|
|
if row is None:
|
|
return
|
|
digest = str(row[0])
|
|
self.object_path(digest).unlink(missing_ok=True)
|
|
connection.execute("DELETE FROM objects WHERE sha256 = ?", (digest,))
|
|
|
|
|
|
class ArtifactGateway:
|
|
"""Resolve named central manifests through a verified local working set."""
|
|
|
|
def __init__(self, store: CentralArtifactStore, cache: LocalArtifactCache) -> None:
|
|
self.store = store
|
|
self.cache = cache
|
|
|
|
def status(self) -> ArtifactGatewayStatus:
|
|
"""Return a lightweight readiness snapshot without resolving an artifact."""
|
|
|
|
try:
|
|
self.store._require_root()
|
|
except ArtifactStoreUnavailable:
|
|
central_status: Literal["ready", "unavailable", "invalid"] = "unavailable"
|
|
except ArtifactIntegrityError:
|
|
central_status = "invalid"
|
|
else:
|
|
central_status = "ready"
|
|
return ArtifactGatewayStatus(
|
|
central_status=central_status,
|
|
cache=self.cache.status(),
|
|
)
|
|
|
|
def publish(
|
|
self,
|
|
*,
|
|
namespace: str,
|
|
key: str,
|
|
artifact_type: str,
|
|
subject_id: str,
|
|
sources: Sequence[tuple[str, str, Path]],
|
|
metadata: Mapping[str, str] | None = None,
|
|
) -> ArtifactManifest:
|
|
members: list[ArtifactMember] = []
|
|
seen_roles: set[str] = set()
|
|
for role, media_type, source in sources:
|
|
_validate_component(role, "artifact role")
|
|
if role in seen_roles:
|
|
raise ValueError(f"duplicate artifact role: {role}")
|
|
seen_roles.add(role)
|
|
published = self.store.publish_file(source)
|
|
members.append(
|
|
ArtifactMember(
|
|
role=role,
|
|
media_type=media_type,
|
|
sha256=published.sha256,
|
|
byte_length=published.byte_length,
|
|
)
|
|
)
|
|
manifest = self.store.publish_manifest(
|
|
artifact_type=artifact_type,
|
|
subject_id=subject_id,
|
|
members=members,
|
|
metadata=metadata,
|
|
)
|
|
self.store.set_reference(namespace, key, manifest.manifest_id)
|
|
reference = self.store.read_reference_document(namespace, key)
|
|
self.cache.save_snapshot(namespace, key, reference, manifest)
|
|
return manifest
|
|
|
|
def resolve_role(self, namespace: str, key: str, role: str) -> ResolvedArtifact:
|
|
manifest, central_available = self._resolve_manifest(namespace, key)
|
|
member = manifest.member(role)
|
|
cached = self.cache.get(member)
|
|
if cached is not None:
|
|
return ResolvedArtifact(
|
|
manifest=manifest,
|
|
member=member,
|
|
path=cached,
|
|
cache_hit=True,
|
|
central_available=central_available,
|
|
)
|
|
if not central_available:
|
|
raise ArtifactStoreUnavailable(
|
|
f"artifact role is not in the offline cache: {namespace}/{key}/{role}"
|
|
)
|
|
path, cache_hit = self.cache.fetch(self.store, member)
|
|
return ResolvedArtifact(
|
|
manifest=manifest,
|
|
member=member,
|
|
path=path,
|
|
cache_hit=cache_hit,
|
|
central_available=True,
|
|
)
|
|
|
|
def pin_reference(self, namespace: str, key: str, *, pin_id: str) -> ArtifactManifest:
|
|
manifest, central_available = self._resolve_manifest(namespace, key)
|
|
fetched: list[ArtifactMember] = []
|
|
try:
|
|
for member in manifest.members:
|
|
if self.cache.get(member) is None:
|
|
if not central_available:
|
|
raise ArtifactStoreUnavailable(
|
|
f"artifact role is not in the offline cache: {member.role}"
|
|
)
|
|
self.cache.fetch(self.store, member, pin_id=pin_id)
|
|
else:
|
|
with self.cache._locked():
|
|
self.cache._pin_locked(pin_id, member.sha256)
|
|
fetched.append(member)
|
|
except BaseException:
|
|
self.cache.unpin(pin_id)
|
|
raise
|
|
return manifest
|
|
|
|
def _resolve_manifest(
|
|
self,
|
|
namespace: str,
|
|
key: str,
|
|
) -> tuple[ArtifactManifest, bool]:
|
|
try:
|
|
reference = self.store.read_reference_document(namespace, key)
|
|
manifest = self.store.read_manifest(str(reference["manifest_id"]))
|
|
except ArtifactStoreUnavailable:
|
|
return self.cache.load_snapshot(namespace, key), False
|
|
self.cache.save_snapshot(namespace, key, reference, manifest)
|
|
return manifest, True
|
|
|
|
|
|
def configured_artifact_gateway(data_dir: Path) -> ArtifactGateway | None:
|
|
"""Build the optional production gateway from environment-only configuration."""
|
|
|
|
configured_root = os.environ.get(STORE_ROOT_ENV, "").strip()
|
|
if not configured_root:
|
|
return None
|
|
max_bytes = _configured_non_negative_int(CACHE_MAX_BYTES_ENV, DEFAULT_CACHE_MAX_BYTES)
|
|
reserve_bytes = _configured_non_negative_int(
|
|
CACHE_RESERVE_BYTES_ENV,
|
|
DEFAULT_FREE_SPACE_RESERVE_BYTES,
|
|
)
|
|
if max_bytes <= 0:
|
|
raise ArtifactCacheCapacityError("artifact cache maximum bytes must be positive")
|
|
return ArtifactGateway(
|
|
CentralArtifactStore(Path(configured_root), create=False),
|
|
LocalArtifactCache(
|
|
data_dir / "artifact-cache",
|
|
max_bytes=max_bytes,
|
|
free_space_reserve_bytes=reserve_bytes,
|
|
),
|
|
)
|
|
|
|
|
|
def _configured_non_negative_int(name: str, default: int) -> int:
|
|
raw = os.environ.get(name, "").strip()
|
|
if not raw:
|
|
return default
|
|
if not raw.isascii() or not raw.isdecimal():
|
|
raise ValueError(f"{name} must be a non-negative integer")
|
|
return int(raw)
|
|
|
|
|
|
def _regular_source(path: Path) -> Path:
|
|
try:
|
|
metadata = path.lstat()
|
|
resolved = path.expanduser().resolve(strict=True)
|
|
except OSError as exc:
|
|
raise ArtifactNotFound("artifact source is unavailable") from exc
|
|
if stat.S_ISLNK(metadata.st_mode) or not stat.S_ISREG(metadata.st_mode):
|
|
raise ArtifactIntegrityError("artifact source must be a regular file")
|
|
return resolved
|
|
|
|
|
|
def _hash_file(path: Path) -> tuple[str, int]:
|
|
digest = hashlib.sha256()
|
|
byte_length = 0
|
|
try:
|
|
with path.open("rb") as stream:
|
|
while chunk := stream.read(_COPY_CHUNK_BYTES):
|
|
digest.update(chunk)
|
|
byte_length += len(chunk)
|
|
except OSError as exc:
|
|
raise ArtifactStoreUnavailable("artifact object is unavailable") from exc
|
|
return digest.hexdigest(), byte_length
|
|
|
|
|
|
def _copy_and_hash(source: Path, destination: Path) -> tuple[str, int]:
|
|
digest = hashlib.sha256()
|
|
byte_length = 0
|
|
with source.open("rb") as input_stream, destination.open("xb") as output_stream:
|
|
while chunk := input_stream.read(_COPY_CHUNK_BYTES):
|
|
output_stream.write(chunk)
|
|
digest.update(chunk)
|
|
byte_length += len(chunk)
|
|
output_stream.flush()
|
|
os.fsync(output_stream.fileno())
|
|
return digest.hexdigest(), byte_length
|
|
|
|
|
|
def _verify_object(path: Path, digest: str, byte_length: int) -> None:
|
|
metadata = path.lstat()
|
|
if stat.S_ISLNK(metadata.st_mode) or not stat.S_ISREG(metadata.st_mode):
|
|
raise ArtifactIntegrityError("content-addressed object path is invalid")
|
|
actual_digest, actual_bytes = _hash_file(path)
|
|
if actual_digest != digest or actual_bytes != byte_length:
|
|
raise ArtifactIntegrityError("content-addressed object is corrupt")
|
|
|
|
|
|
def _validate_members(members: Sequence[ArtifactMember]) -> None:
|
|
if not members:
|
|
raise ValueError("artifact manifest must contain at least one member")
|
|
roles: set[str] = set()
|
|
for member in members:
|
|
_validate_member(member)
|
|
if member.role in roles:
|
|
raise ValueError(f"duplicate artifact role: {member.role}")
|
|
roles.add(member.role)
|
|
|
|
|
|
def _validate_member(member: ArtifactMember) -> None:
|
|
_validate_component(member.role, "artifact role")
|
|
if not isinstance(member.media_type, str) or not 1 <= len(member.media_type) <= 255:
|
|
raise ValueError("artifact media type is invalid")
|
|
_validate_sha256(member.sha256)
|
|
if not isinstance(member.byte_length, int) or isinstance(member.byte_length, bool):
|
|
raise ValueError("artifact byte length is invalid")
|
|
if member.byte_length < 0:
|
|
raise ValueError("artifact byte length is invalid")
|
|
|
|
|
|
def _parse_manifest(document: Mapping[str, Any], manifest_id: str) -> ArtifactManifest:
|
|
if (
|
|
document.get("schema_version") != MANIFEST_SCHEMA
|
|
or not isinstance(document.get("artifact_type"), str)
|
|
or not isinstance(document.get("subject_id"), str)
|
|
or not isinstance(document.get("created_at_utc"), str)
|
|
or not isinstance(document.get("members"), list)
|
|
or not isinstance(document.get("metadata"), dict)
|
|
):
|
|
raise ArtifactIntegrityError("artifact manifest shape is invalid")
|
|
_validate_component(str(document["artifact_type"]), "artifact type")
|
|
_validate_component(str(document["subject_id"]), "artifact subject")
|
|
raw_metadata = document["metadata"]
|
|
assert isinstance(raw_metadata, dict)
|
|
metadata: dict[str, str] = {}
|
|
for key, value in raw_metadata.items():
|
|
if not isinstance(key, str) or not isinstance(value, str):
|
|
raise ArtifactIntegrityError("artifact manifest metadata is invalid")
|
|
_validate_component(key, "metadata key")
|
|
metadata[key] = value
|
|
raw_members = document["members"]
|
|
assert isinstance(raw_members, list)
|
|
members: list[ArtifactMember] = []
|
|
for raw in raw_members:
|
|
if not isinstance(raw, dict):
|
|
raise ArtifactIntegrityError("artifact manifest member is invalid")
|
|
try:
|
|
member = ArtifactMember(
|
|
role=str(raw["role"]),
|
|
media_type=str(raw["media_type"]),
|
|
sha256=str(raw["sha256"]),
|
|
byte_length=int(raw["byte_length"]),
|
|
)
|
|
except (KeyError, TypeError, ValueError) as exc:
|
|
raise ArtifactIntegrityError("artifact manifest member is invalid") from exc
|
|
_validate_member(member)
|
|
members.append(member)
|
|
_validate_members(members)
|
|
if tuple(members) != tuple(sorted(members, key=lambda item: item.role)):
|
|
raise ArtifactIntegrityError("artifact manifest members are not canonical")
|
|
return ArtifactManifest(
|
|
manifest_id=manifest_id,
|
|
artifact_type=str(document["artifact_type"]),
|
|
subject_id=str(document["subject_id"]),
|
|
created_at_utc=str(document["created_at_utc"]),
|
|
members=tuple(members),
|
|
metadata=metadata,
|
|
)
|
|
|
|
|
|
def _manifest_document(manifest: ArtifactManifest) -> dict[str, Any]:
|
|
return {
|
|
"schema_version": MANIFEST_SCHEMA,
|
|
"artifact_type": manifest.artifact_type,
|
|
"subject_id": manifest.subject_id,
|
|
"created_at_utc": manifest.created_at_utc,
|
|
"members": [
|
|
{
|
|
"role": item.role,
|
|
"media_type": item.media_type,
|
|
"sha256": item.sha256,
|
|
"byte_length": item.byte_length,
|
|
}
|
|
for item in manifest.members
|
|
],
|
|
"metadata": dict(sorted(manifest.metadata.items())),
|
|
}
|
|
|
|
|
|
def _validate_reference_snapshot(
|
|
document: Mapping[str, Any],
|
|
namespace: str,
|
|
key: str,
|
|
) -> None:
|
|
if (
|
|
document.get("schema_version") != REFERENCE_SCHEMA
|
|
or document.get("namespace") != namespace
|
|
or document.get("key") != key
|
|
or not isinstance(document.get("manifest_id"), str)
|
|
or _SHA256.fullmatch(str(document["manifest_id"])) is None
|
|
or not isinstance(document.get("updated_at_utc"), str)
|
|
):
|
|
raise ArtifactIntegrityError("artifact reference snapshot is invalid")
|
|
|
|
|
|
def _read_json_document(path: Path, *, unavailable_message: str) -> dict[str, Any]:
|
|
try:
|
|
metadata = path.lstat()
|
|
if (
|
|
stat.S_ISLNK(metadata.st_mode)
|
|
or not stat.S_ISREG(metadata.st_mode)
|
|
or not 0 < metadata.st_size <= _MAX_MANIFEST_BYTES
|
|
):
|
|
raise ArtifactIntegrityError("artifact metadata path is invalid")
|
|
value = json.loads(path.read_text(encoding="utf-8"))
|
|
except FileNotFoundError as exc:
|
|
raise ArtifactStoreUnavailable(unavailable_message) from exc
|
|
except OSError as exc:
|
|
raise ArtifactStoreUnavailable(unavailable_message) from exc
|
|
except json.JSONDecodeError as exc:
|
|
raise ArtifactIntegrityError("artifact metadata is not valid JSON") from exc
|
|
if not isinstance(value, dict):
|
|
raise ArtifactIntegrityError("artifact metadata is not an object")
|
|
return value
|
|
|
|
|
|
def _write_bytes_atomic(path: Path, payload: bytes, *, mode: int | None = None) -> None:
|
|
path.parent.mkdir(parents=True, exist_ok=True)
|
|
temporary = path.with_name(f".{path.name}.{uuid4().hex}.tmp")
|
|
try:
|
|
with temporary.open("xb") as stream:
|
|
stream.write(payload)
|
|
stream.flush()
|
|
os.fsync(stream.fileno())
|
|
if mode is not None:
|
|
os.chmod(temporary, mode)
|
|
os.replace(temporary, path)
|
|
finally:
|
|
temporary.unlink(missing_ok=True)
|
|
|
|
|
|
def _canonical_json(value: object) -> bytes:
|
|
return json.dumps(
|
|
value,
|
|
ensure_ascii=False,
|
|
sort_keys=True,
|
|
separators=(",", ":"),
|
|
allow_nan=False,
|
|
).encode("utf-8")
|
|
|
|
|
|
def _validate_component(value: str, label: str) -> None:
|
|
if not isinstance(value, str) or _SAFE_COMPONENT.fullmatch(value) is None:
|
|
raise ValueError(f"{label} is invalid")
|
|
|
|
|
|
def _validate_sha256(value: str) -> None:
|
|
if not isinstance(value, str) or _SHA256.fullmatch(value) is None:
|
|
raise ValueError("artifact SHA-256 is invalid")
|