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
51 changes: 43 additions & 8 deletions build_system/builder/cache/inventory.py
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,7 @@
from .cargounits import unaccounted_size, unit_entries
from .contract import PruneStrategy
from .inventorymodels import RetentionInventory
from .leases import active_path
from .leases import active_path, lease_key
from .measure import measure
from .models import CacheEntry, CacheInventory, CachePolicy, StageInventory
from .objectunits import object_entries
Expand All @@ -25,6 +25,23 @@ def _lease_active(stage_path: Path, template: str | None, key: str) -> bool:
return active_path(lease)


def _managed(stage_policy, name: str) -> bool:
return any(fnmatch.fnmatchcase(name, pattern) for pattern in stage_policy.managed_globs)


def _lease_files(template: str | None, children: list[Path], stage_policy) -> dict[str, str]:
"""Map each lease file's name to the generation key it holds."""
if template is None:
return {}
found = {}
for child in children:
key = lease_key(template, child.name)
if key is not None and not _managed(stage_policy, child.name) and not child.is_symlink() \
and child.is_file():
found[child.name] = key
return found


def _entry_size(path: Path, allocated_seen: set[tuple[int, int]]) -> tuple[int, int]:
measured = measure(path, allocated_seen)
return measured.logical_bytes, measured.allocated_bytes
Expand Down Expand Up @@ -54,25 +71,43 @@ def _stage_inventory(
unmanaged_allocated = 0
busy = any(active_path(stage_root / lock) for lock in stage_policy.mutation_locks)
if stage_path.is_dir():
for child in sorted(stage_path.iterdir(), key=lambda item: item.name):
children = sorted(stage_path.iterdir(), key=lambda item: item.name)
names = {child.name for child in children}
managed_names = {name for name in names if _managed(stage_policy, name)}
leases = _lease_files(stage_policy.lease_template, children, stage_policy)
for child in children:
key = leases.get(child.name)
if key is not None and key in managed_names:
continue # removed and accounted with its generation
if key in names:
key = None # the lease of an unmanaged sibling stays unmanaged
logical, allocated = _entry_size(child, allocated_seen)
stat = child.lstat()
managed = any(
fnmatch.fnmatchcase(child.name, pattern) for pattern in stage_policy.managed_globs
)
lease_only = key is not None
managed = lease_only or child.name in managed_names
members: tuple[Path, ...] = ()
if managed and not lease_only and stage_policy.lease_template is not None:
lease = stage_policy.lease_template.format(key=child.name)
if lease in leases:
members = (Path(lease),)
lease_logical, lease_allocated = _entry_size(stage_path / lease, allocated_seen)
logical += lease_logical
allocated += lease_allocated
entries.append(
CacheEntry(
key=child.name,
key=key or child.name,
relative_path=Path(child.name),
member_paths=members,
logical_bytes=logical,
allocated_bytes=allocated,
created_ns=stat.st_ctime_ns,
last_used_ns=stat.st_atime_ns,
managed=managed,
lease_only=lease_only,
protected=managed
and (
busy or child.name in referenced
or _lease_active(stage_path, stage_policy.lease_template, child.name)
busy or (key or child.name) in referenced
or _lease_active(stage_path, stage_policy.lease_template, key or child.name)
),
)
)
Expand Down
62 changes: 56 additions & 6 deletions build_system/builder/cache/leases.py
Original file line number Diff line number Diff line change
Expand Up @@ -26,22 +26,72 @@


def retain_path(path: Path) -> BinaryIO:
"""Hold one shared lease until explicit release or process exit."""
"""Hold one shared lease until explicit release or process exit.

The lock waits out an exclusive holder: a pruner holds one only while it
checks or removes a dead generation. A lease file the pruner unlinked
meanwhile is a lock nobody can see, so it is retaken on the live path.
"""
lease = path.absolute()
existing = _HELD.get(lease)
if existing is not None and not existing.closed:
return existing
lease.parent.mkdir(parents=True, exist_ok=True)
descriptor = os.fdopen(os.open(lease, os.O_APPEND | os.O_CREAT | os.O_RDWR, 0o600), "a+b")
try:
fcntl.flock(descriptor, fcntl.LOCK_SH | fcntl.LOCK_NB)
except BaseException:
while True:
descriptor = os.fdopen(
os.open(lease, os.O_APPEND | os.O_CREAT | os.O_RDWR | os.O_NOFOLLOW, 0o600), "a+b"
)
try:
fcntl.flock(descriptor, fcntl.LOCK_SH)
if _same_file(descriptor, lease):
break
except BaseException:
descriptor.close()
raise
descriptor.close()
raise
_HELD[lease] = descriptor
return descriptor


def _same_file(descriptor: BinaryIO, path: Path) -> bool:
try:
current = path.lstat()
except FileNotFoundError:
return False
held = os.fstat(descriptor.fileno())
return (held.st_dev, held.st_ino) == (current.st_dev, current.st_ino)


def lease_key(template: str, name: str) -> str | None:
"""The generation key a lease file name encodes under `template`, if any."""
prefix, suffix = template.split("{key}")
if len(name) <= len(prefix) + len(suffix):
return None
if not (name.startswith(prefix) and name.endswith(suffix)):
return None
return name[len(prefix):len(name) - len(suffix)]


@contextmanager
def exclusive_path(path: Path) -> Iterator[bool]:
"""Hold a lease exclusively while its generation is removed.

Yields False when an owner holds it. Creating the file when it is missing
is deliberate: an owner leases before it creates its generation, so a
concurrent owner then waits on this inode and retakes the live path.
"""
descriptor = os.fdopen(
os.open(path, os.O_APPEND | os.O_CREAT | os.O_RDWR | os.O_NOFOLLOW, 0o600), "a+b"
)
with descriptor:
try:
fcntl.flock(descriptor, fcntl.LOCK_EX | fcntl.LOCK_NB)
except BlockingIOError:
yield False
return
yield True


def retain_generation(paths: CachePaths, stage_id: str, key: str) -> BinaryIO:
"""Hold the configured lease for one managed stage generation."""
stage = paths.policy.stages[stage_id]
Expand Down
4 changes: 4 additions & 0 deletions build_system/builder/cache/models.py
Original file line number Diff line number Diff line change
Expand Up @@ -160,6 +160,8 @@ class CacheEntry(BaseModel):
last_used_ns: Annotated[StrictInt, Field(ge=0)]
managed: StrictBool = True
protected: StrictBool = False
#: A lease file whose generation is gone: collected, never counted.
lease_only: StrictBool = False


class StageInventory(BaseModel):
Expand Down Expand Up @@ -223,6 +225,8 @@ class ApplyResult(BaseModel):

removed: tuple[Path, ...]
missing: tuple[Path, ...]
#: Planned paths kept because an owner leased their generation meanwhile.
busy: tuple[Path, ...] = ()
journal: Path


Expand Down
60 changes: 50 additions & 10 deletions build_system/builder/cache/operations.py
Original file line number Diff line number Diff line change
Expand Up @@ -8,12 +8,14 @@
import stat
import time
from collections.abc import Iterator
from contextlib import ExitStack
from pathlib import Path

from pydantic import ValidationError

from .leases import mutation_locks
from .models import AdmissionEvent, ApplyResult, PrunePlan
from .leases import exclusive_path, mutation_locks, release_path
from .measure import measure
from .models import AdmissionEvent, ApplyResult, PruneAction, PrunePlan
from .paths import CachePaths

JOURNAL_PATH = Path("state/events/cache.jsonl")
Expand Down Expand Up @@ -72,16 +74,28 @@ def apply_prune(paths: CachePaths, plan: PrunePlan, *, reason: str) -> ApplyResu
if not reason.strip():
raise ValueError("cache mutation reason must be non-empty")
targets = tuple(paths.contained_entry(action.stage_id, action.path) for action in plan.actions)
generations: dict[tuple[str, str], list[Path]] = {}
for action, target in zip(plan.actions, targets, strict=True):
generations.setdefault((action.stage_id, action.key), []).append(target)
removed: list[Path] = []
missing: list[Path] = []
busy: list[Path] = []
with mutation_locks(paths, (action.stage_id for action in plan.actions)) as locks:
for selected in targets:
for target in _unlocked_targets(selected, locks):
if target.exists() or target.is_symlink():
_remove(target)
removed.append(target)
else:
missing.append(target)
for (stage_id, key), selected in generations.items():
template = paths.policy.stages[stage_id].lease_template
root = paths.stage(stage_id)
lease = None if template is None or not root.is_dir() else root / template.format(key=key)
with ExitStack() as stack:
if lease is not None and not stack.enter_context(exclusive_path(lease)):
busy.extend(selected)
continue
for target in (*selected, *(() if lease is None else (lease,))):
for unlocked in _unlocked_targets(target, locks):
if unlocked.exists() or unlocked.is_symlink():
_remove(unlocked)
removed.append(unlocked)
elif unlocked != lease:
missing.append(unlocked)
journal = paths.root / JOURNAL_PATH
journal.parent.mkdir(parents=True, exist_ok=True)
event = {
Expand All @@ -91,11 +105,37 @@ def apply_prune(paths: CachePaths, plan: PrunePlan, *, reason: str) -> ApplyResu
"reason": reason,
"removed": [str(path) for path in removed],
"missing": [str(path) for path in missing],
"busy": [str(path) for path in busy],
"reclaim_bytes": plan.reclaim_bytes,
}
with journal.open("a", encoding="utf-8") as stream:
stream.write(json.dumps(event, sort_keys=True) + "\n")
return ApplyResult(removed=tuple(removed), missing=tuple(missing), journal=journal)
return ApplyResult(
removed=tuple(removed), missing=tuple(missing), busy=tuple(busy), journal=journal
)


def reclaim_generation(paths: CachePaths, stage_id: str, key: str, *, reason: str) -> ApplyResult:
"""End one generation this process leased: release the lease, then remove
the generation and its lease through the same guarded, journaled path a
prune takes. A concurrent prune that wins the lease first does the same
removal; whichever loses sees it busy or already gone."""
stage = paths.policy.stages[stage_id]
if stage.lease_template is None:
raise ValueError(f"cache stage {stage_id!r} has no generation lease")
root = paths.stage(stage_id)
generation = root / key
release_path(root / stage.lease_template.format(key=key))
logical = measure(generation, set()).logical_bytes if generation.exists() else 0
plan = PrunePlan(
generated_ns=time.time_ns(),
reclaim_bytes=logical,
actions=(PruneAction(
stage_id=stage_id, key=key, path=generation, logical_bytes=logical, reason=reason,
),),
violations=(),
)
return apply_prune(paths, plan, reason=reason)


def record_admission_event(root: Path, state_path: Path, event: AdmissionEvent) -> Path:
Expand Down
54 changes: 33 additions & 21 deletions build_system/builder/cache/planner.py
Original file line number Diff line number Diff line change
Expand Up @@ -20,31 +20,35 @@ def _age_clock(strategy: PruneStrategy, entry: CacheEntry) -> int:
return entry.last_used_ns if strategy is PruneStrategy.LRU else entry.created_ns


def _actions(stage, entry: CacheEntry, reason: str) -> tuple[PruneAction, ...]:
"""A generation with members is removed whole and in order; its bytes are
carried once, on the first path."""
return tuple(
PruneAction(
stage_id=stage.stage_id,
key=entry.key,
path=stage.path / relative,
logical_bytes=entry.logical_bytes if index == 0 else 0,
reason=reason,
)
for index, relative in enumerate((entry.relative_path, *entry.member_paths))
)


def plan_prune(inventory: CacheInventory | RetentionInventory, policy: CachePolicy) -> PrunePlan:
"""Select expired, surplus, and pressure candidates without touching pinned state."""
actions: list[PruneAction] = []
violations: list[str] = []
selected: dict[str, set[str]] = {stage.stage_id: set() for stage in inventory.stages}

def choose(stage, entry, reason: str) -> None:
# A generation with members is removed whole and in order; its bytes
# are carried once, on the first path.
for index, relative in enumerate((entry.relative_path, *entry.member_paths)):
actions.append(
PruneAction(
stage_id=stage.stage_id,
key=entry.key,
path=stage.path / relative,
logical_bytes=entry.logical_bytes if index == 0 else 0,
reason=reason,
)
)
actions.extend(_actions(stage, entry, reason))
selected[stage.stage_id].add(entry.key)

for stage in inventory.stages:
stage_policy = policy.stages[stage.stage_id]
remaining = stage.logical_bytes
remaining_count = sum(entry.managed for entry in stage.entries)
remaining_count = sum(entry.managed and not entry.lease_only for entry in stage.entries)
if stage_policy.prune_strategy is PruneStrategy.NONE:
if remaining > stage_policy.max_size_bytes:
violations.append(
Expand All @@ -66,7 +70,20 @@ def choose(stage, entry, reason: str) -> None:
)
maximum_age = stage_policy.maximum_age_hours * NANOSECONDS_PER_HOUR
over_max = remaining > stage_policy.max_size_bytes
# A leased ephemeral generation lives exactly as long as its owner:
# once no process holds the lease, it is garbage whatever its age.
ownerless = (
stage_policy.prune_strategy is PruneStrategy.EPHEMERAL
and stage_policy.lease_template is not None
)
for entry in ordered:
if entry.protected:
continue
if entry.lease_only or ownerless:
choose(stage, entry, "orphaned lease" if entry.lease_only else "no live owner")
remaining -= entry.logical_bytes
remaining_count -= not entry.lease_only
continue
expired = (
inventory.generated_ns
>= _age_clock(stage_policy.prune_strategy, entry) + maximum_age
Expand All @@ -76,7 +93,7 @@ def choose(stage, entry, reason: str) -> None:
stage_policy.maximum_count is not None
and remaining_count > stage_policy.maximum_count
)
if entry.protected or not (expired or recover_to_warm or over_count):
if not (expired or recover_to_warm or over_count):
continue
reason = (
"expired"
Expand Down Expand Up @@ -117,16 +134,11 @@ def plan_clean(inventory: CacheInventory, stage_id: str) -> PrunePlan:
known = ", ".join(stage.stage_id for stage in inventory.stages)
raise KeyError(f"unknown cache stage {stage_id!r}; expected one of: {known}, all")
actions = tuple(
PruneAction(
stage_id=stage.stage_id,
key=entry.key,
path=stage.path / entry.relative_path,
logical_bytes=entry.logical_bytes,
reason="explicit clean",
)
action
for stage in selected
for entry in stage.entries
if entry.managed and not entry.protected
for action in _actions(stage, entry, "explicit clean")
)
return PrunePlan(
generated_ns=inventory.generated_ns,
Expand Down
Loading
Loading