Coverage for cuda/core/utils/_program_cache/_file_stream.py: 89.08%
284 statements
« prev ^ index » next coverage.py v7.15.2, created at 2026-07-29 01:38 +0000
« prev ^ index » next coverage.py v7.15.2, created at 2026-07-29 01:38 +0000
1# SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved.
2#
3# SPDX-License-Identifier: Apache-2.0
5"""On-disk bytes-in / bytes-out program cache.
7Atomic writes via :func:`os.replace`. Concurrent readers see either the
8old entry or the new one, never a partial file. Each entry is the raw
9compiled binary so files are directly consumable by external NVIDIA
10tools (``cuobjdump``, ``nvdisasm``, ``cuda-gdb``).
11"""
13from __future__ import annotations
15import contextlib
16import errno
17import hashlib
18import os
19import tempfile
20import threading
21import time
22from pathlib import Path
23from typing import Any, Callable, Iterable
25from cuda.core._module import ObjectCode
27from ._abc import ProgramCacheResource, _as_key_bytes, _extract_bytes
29_ENTRIES_SUBDIR = "entries"
30_TMP_SUBDIR = "tmp"
31# Temp files older than this are assumed to belong to a crashed writer and
32# are eligible for cleanup. Picked large enough that no real ``os.replace``
33# write should still be in flight (writes are bounded by mkstemp + write +
34# fsync + replace, all fast on healthy disks).
35_TMP_STALE_AGE_SECONDS = 3600
38_SHARING_VIOLATION_WINERRORS = (5, 32, 33) # ERROR_ACCESS_DENIED, ERROR_SHARING_VIOLATION, ERROR_LOCK_VIOLATION
39_REPLACE_RETRY_DELAYS = (0.0, 0.005, 0.010, 0.020, 0.050, 0.100) # ~185ms budget
42# Exposed as a module-level flag so tests can toggle it without monkeypatching
43# ``os.name`` itself (pathlib reads ``os.name`` at instantiation time).
44_IS_WINDOWS = os.name == "nt"
47def _stat_key(st: os.stat_result) -> tuple[int, int, int]:
48 """Stat fingerprint used by every stat-guarded path.
50 ``(st_ino, st_size, st_mtime_ns)`` is the smallest triple that
51 distinguishes "same file" from "file replaced under us": ``st_ino``
52 catches replacement, ``st_size`` and ``st_mtime_ns`` catch a write
53 that happens to land on the same inode (e.g. truncate-and-write in
54 place). Centralised so all four readers compare the same fields.
55 """
56 return (st.st_ino, st.st_size, st.st_mtime_ns) 1rvwxybopmqfezAnBWClDEacPgdFhjGM23st
59def _default_cache_dir() -> Path:
60 """OS-conventional default location for the file-stream cache.
62 Resolves to the user-cache root for the calling user, with a
63 ``program-cache`` leaf so future tooling can place sibling caches
64 under the same ``cuda-python`` vendor directory:
66 * Linux: ``$XDG_CACHE_HOME/cuda-python/program-cache``
67 (default ``~/.cache/cuda-python/program-cache`` per the XDG Base
68 Directory spec).
69 * Windows: ``%LOCALAPPDATA%\\cuda-python\\program-cache``
70 (Windows uses local AppData -- caches don't roam; falls back to
71 ``~/AppData/Local`` if the env var is unset).
73 CUDA does not support macOS, so no macOS branch is provided.
74 """
75 if _IS_WINDOWS: 1)
76 local_app_data = os.environ.get("LOCALAPPDATA") 1)
77 root = Path(local_app_data) if local_app_data else Path.home() / "AppData" / "Local" 1)
78 else:
79 xdg = os.environ.get("XDG_CACHE_HOME") 1)
80 root = Path(xdg) if xdg else Path.home() / ".cache" 1)
81 return root / "cuda-python" / "program-cache" 1)
84def _with_sharing_retry(
85 op: Callable[..., Any], *args: Any, on_exhausted: Callable[..., Any] | None = None, **kwargs: Any
86) -> Any:
87 """Run ``op(*args, **kwargs)`` retrying transient Windows sharing
88 violations under the bounded ``_REPLACE_RETRY_DELAYS`` budget.
90 On Windows, ``os.replace``/``read_bytes``/``unlink`` can surface
91 winerror 5/32/33 (or bare EACCES via ``_is_windows_sharing_violation``)
92 while another process briefly holds the file open without share-delete
93 rights. The retry hides that contention. Other ``PermissionError``s
94 (real ACLs, unexpected winerror) propagate immediately.
96 Successful returns and any non-``PermissionError`` exceptions
97 (including ``FileNotFoundError``) bubble up unchanged. After the
98 budget is exhausted, the helper either calls ``on_exhausted(last_exc)``
99 if provided, or re-raises the last sharing-violation exception.
100 """
101 last_exc: PermissionError | None = None 1rvwxybRopmqK6fezAXnZHITUJBWClDEakcPigYQudFNhjS1V0LOGM423st
102 for delay in _REPLACE_RETRY_DELAYS: 1rvwxybRopmqK6fezAXnZHITUJBWClDEakcPigYQudFNhjS1V0LOGM423st
103 if delay: 1rvwxybRopmqK6fezAXnZHITUJBWClDEakcPigYQudFNhjS1V0LOGM423st
104 time.sleep(delay) 1HIJlud
105 try: 1rvwxybRopmqK6fezAXnZHITUJBWClDEakcPigYQudFNhjS1V0LOGM423st
106 return op(*args, **kwargs) 1rvwxybRopmqK6fezAXnZHITUJBWClDEakcPigYQudFNhjS1V0LOGM423st
107 except PermissionError as exc: 1rbRK6feZHITUJlacQudst
108 if not _is_windows_sharing_violation(exc): 1feZHITUJlQud
109 raise 1feZTUQ
110 last_exc = exc 1HIJlud
111 if on_exhausted is not None: 1HIJd
112 return on_exhausted(last_exc) 1HIJ
113 assert last_exc is not None # at least one iteration ran and caught a PermissionError 1d
114 raise last_exc 1d
117def _replace_with_sharing_retry(tmp_path: Path, target: Path) -> bool:
118 """Atomic rename with Windows-specific retry on sharing/lock violations.
120 Returns True on success. Returns False only after the retry budget is
121 exhausted on Windows with a genuine sharing violation -- the caller then
122 treats the cache write as dropped. Any other ``PermissionError`` (ACLs,
123 read-only dir, unexpected winerror, or any POSIX failure) propagates.
125 ``ERROR_ACCESS_DENIED`` (winerror 5) is treated as a sharing violation
126 because Windows surfaces it when a file is held open without
127 ``FILE_SHARE_WRITE`` (Python's default for ``open(p, "wb")``) or while
128 a previous unlink is in ``PENDING_DELETE`` -- both are transient.
129 """
131 def _do_replace() -> bool: 1rvwxybRopmqKfezAXnZHITUJBWClDEakcPigYQudFNhjS1V0LOG4st
132 os.replace(tmp_path, target) 1rvwxybRopmqKfezAXnZHITUJBWClDEakcPigYQudFNhjS1V0LOG4st
133 return True 1rvwxybopmqKfezAXnBWClDEakcPigYQudFNhjS1V0LOG4st
135 return bool(_with_sharing_retry(_do_replace, on_exhausted=lambda _exc: False)) 1rvwxybRopmqKfezAXnZHITUJBWClDEakcPigYQudFNhjS1V0LOG4st
138def _stat_and_read_with_sharing_retry(path: Path) -> tuple[os.stat_result, bytes]:
139 """Snapshot stat and read bytes, retrying briefly on Windows transient
140 sharing-violation ``PermissionError``.
142 Reads race the rewriter's ``os.replace``: on Windows, the destination
143 can be momentarily inaccessible (winerror 5/32/33) while the rename
144 completes. Mirroring ``_replace_with_sharing_retry``'s budget keeps
145 transient contention from being mistaken for a real read failure.
147 Raises ``FileNotFoundError`` on miss or after exhausting the Windows
148 sharing-retry budget. Non-Windows ``PermissionError`` propagates.
150 On Windows, EACCES (errno 13) is treated as transient too: ``io.open``
151 sometimes surfaces a pending-delete or share-mode mismatch as bare
152 EACCES with no ``winerror`` attribute, indistinguishable here from
153 a true sharing violation. Real ACL problems on a path the cache owns
154 would surface consistently; the bounded retry budget keeps the cost
155 of treating them as transient negligible.
156 """
158 def _do_stat_and_read() -> tuple[os.stat_result, bytes]: 1rvwxybRK6zAnHIJBClDEacuFLOGM23st
159 return path.stat(), path.read_bytes() 1rvwxybRK6zAnHIJBClDEacuFLOGM23st
161 def _exhausted(last_exc: PermissionError) -> None: 1rvwxybRK6zAnHIJBClDEacuFLOGM23st
162 raise FileNotFoundError(path) from last_exc
164 return _with_sharing_retry(_do_stat_and_read, on_exhausted=_exhausted) # type: ignore[no-any-return] 1rvwxybRK6zAnHIJBClDEacuFLOGM23st
167_UTIME_SUPPORTS_FD = os.utime in os.supports_fd
170def _touch_atime(path: Path, st_before: os.stat_result) -> None:
171 """Bump ``path``'s atime to "now", preserving its mtime, iff the
172 file's stat still matches ``st_before``.
174 Eviction sorts by ``st_atime`` so reads must reliably refresh atime
175 regardless of OS or filesystem default behavior:
177 * Linux ``relatime`` (default) only updates atime when the existing
178 atime is older than mtime, which would skew LRU once an entry has
179 been read once.
180 * NTFS on Windows Vista+ disables atime updates by default
181 (``NtfsDisableLastAccessUpdate``) and most modern installations
182 keep that off, so a bare read never bumps atime.
183 * ``noatime``-mounted filesystems disable updates entirely.
185 Calling ``os.utime`` with explicit times bypasses all of the above
186 and writes atime directly. The stat-guard is critical: if another
187 process ``os.replace``-d a fresh entry into ``path`` between the
188 read and this touch, blindly applying ``st_before.st_mtime_ns``
189 would roll the new entry's mtime back to the old value and confuse
190 the eviction stat-guard (which checks ``(ino, size, mtime_ns)``)
191 into deleting a freshly-committed file.
193 Where ``os.utime`` supports file descriptors (Linux, macOS), the
194 fstat-then-utime pair runs against the same open fd: even if another
195 writer replaces the path between our ``os.open`` and the ``fstat``,
196 the fd still refers to the file we opened, so the comparison and the
197 utime both target the same inode. This closes the residual TOCTOU
198 window that a path-based stat + path-based utime would have.
200 On Windows, ``os.utime`` is path-only; the fallback re-stats the
201 path and accepts a small TOCTOU window between the second stat and
202 the utime. That window is microseconds and the worst-case outcome
203 is the racing writer's mtime being rolled back by a few hundred
204 nanoseconds -- the eviction stat-guard would then refuse to evict
205 the slightly-stale entry, costing one cache miss (recompile) but
206 not a corrupt eviction.
208 Best-effort: any ``OSError`` (read-only mount, restrictive ACLs,
209 ...) is swallowed -- size enforcement still bounds the cache, but
210 eviction degrades toward FIFO.
211 """
212 new_atime_ns = time.time_ns() 1rvwxybzAnBClDEacPF0LOGM23st
213 if _UTIME_SUPPORTS_FD: 1rvwxybzAnBClDEacPF0LOGM23st
214 try: 1rvwxybzAnBClDEacPFLOGM23st
215 fd = os.open(path, os.O_RDONLY) 1rvwxybzAnBClDEacPFLOGM23st
216 except OSError: 1O
217 return 1O
218 try: 1rvwxybzAnBClDEacPFLGM23st
219 try: 1rvwxybzAnBClDEacPFLGM23st
220 st_now = os.fstat(fd) 1rvwxybzAnBClDEacPFLGM23st
221 except OSError: 1L
222 return 1L
223 if _stat_key(st_now) != _stat_key(st_before): 1rvwxybzAnBClDEacPFGM23st
224 return 1P
225 with contextlib.suppress(OSError): 1rvwxybzAnBClDEacPFGM23st
226 os.utime(fd, ns=(new_atime_ns, st_before.st_mtime_ns)) 1rvwxybzAnBClDEacPFGM23st
227 finally:
228 os.close(fd) 1rvwxybzAnBClDEacPFLGM23st
229 return 1rvwxybzAnBClDEacPFGM23st
231 # Path-based fallback (Windows). Best-effort -- residual TOCTOU window
232 # documented above.
233 try: 10
234 st_now = path.stat() 10
235 except OSError: 10
236 return 10
237 if _stat_key(st_now) != _stat_key(st_before):
238 return
239 with contextlib.suppress(OSError):
240 os.utime(path, ns=(new_atime_ns, st_before.st_mtime_ns))
243def _is_windows_sharing_violation(exc: BaseException) -> bool:
244 """Return True if ``exc`` is a Windows sharing/lock violation that
245 :func:`_unlink_with_sharing_retry` would have retried.
247 Used by best-effort callers to filter out the exhausted-retry case
248 while letting other ``PermissionError`` instances (POSIX ACL
249 issues, Windows non-sharing winerrors) propagate -- those are real
250 configuration problems, not transient contention.
252 The ``EACCES`` fallback only fires when ``winerror`` is absent: a
253 bare ``EACCES`` (no winerror attached) is the way ``io.open``
254 surfaces a pending-delete or share-mode mismatch on Windows. When
255 ``winerror`` IS set but is NOT in the sharing set, the OS told us
256 exactly what failed and it isn't a sharing violation -- treating it
257 as transient would silently swallow real errors like a corrupt
258 ACL.
259 """
260 if not _IS_WINDOWS: 1feZHITUJlQud(
261 return False 1fZQ(
262 if not isinstance(exc, PermissionError): 1eHITUJlud(
263 return False
264 winerror = getattr(exc, "winerror", None) 1eHITUJlud(
265 if winerror in _SHARING_VIOLATION_WINERRORS: 1eHITUJlud(
266 return True 1HIJlud(
267 return winerror is None and exc.errno == errno.EACCES 1eTU(
270def _unlink_with_sharing_retry(path: Path) -> None:
271 """Unlink with Windows-specific retry on sharing/lock violations.
273 On Windows, ``Path.unlink`` raises ``PermissionError`` (winerror 5,
274 32, or 33; sometimes bare ``EACCES``) when another process holds
275 the file open without ``FILE_SHARE_DELETE``. Python's default
276 ``open(p, "rb")`` does not pass that flag, so a reader from another
277 process briefly blocks our unlink while it reads. Retry with the
278 same backoff budget as :func:`_replace_with_sharing_retry` so
279 transient contention is not turned into a propagated error.
281 Raises ``FileNotFoundError`` if the file is absent; the last
282 ``PermissionError`` if the Windows retry budget is exhausted; and
283 propagates any non-sharing ``PermissionError`` (or any non-Windows
284 ``PermissionError``) immediately. Best-effort callers should use
285 :func:`_is_windows_sharing_violation` to filter the exhausted-retry
286 case and re-raise any other ``PermissionError``.
287 """
288 _with_sharing_retry(path.unlink) 1bopmqKfeWacigQudhj
291def _prune_if_stat_unchanged(path: Path, st_before: os.stat_result) -> None:
292 """Unlink ``path`` iff its stat still matches ``st_before``.
294 Guards against a cross-process race: a reader that sees a corrupt
295 record can have it atomically replaced (via ``os.replace``) by a
296 writer before the reader decides to prune. Comparing
297 ``(ino, size, mtime_ns)`` before and after rules out that case --
298 any mismatch means someone else wrote a new file and we must not
299 delete their work. The residual TOCTOU window between stat and
300 unlink is narrow; worst case, a very-recently-written entry is
301 removed and the next read recompiles.
303 Best-effort: a Windows sharing violation that survives the retry
304 budget leaves the file in place. The caller is in an eviction or
305 cleanup pass, so re-trying on the next pass is the right outcome.
306 """
307 try: 1opmqWj
308 st_now = path.stat() 1opmqWj
309 except FileNotFoundError:
310 return
311 if _stat_key(st_before) != _stat_key(st_now): 1opmqWj
312 return 1pW
313 try: 1opmqWj
314 _unlink_with_sharing_retry(path) 1opmqWj
315 except FileNotFoundError:
316 pass
317 except PermissionError as exc:
318 # Swallow only the exhausted-Windows-sharing case. POSIX ACL
319 # errors and Windows non-sharing winerrors are real configuration
320 # problems and must surface, not be silently lost during a prune.
321 if not _is_windows_sharing_violation(exc):
322 raise
325class FileStreamProgramCache(ProgramCacheResource):
326 """Persistent program cache backed by a directory of atomic files.
328 Designed for multi-process use: writes stage a temporary file and then
329 :func:`os.replace` it into place, so concurrent readers never observe a
330 partially-written entry. Each entry on disk is the raw compiled binary
331 -- cubin / PTX / LTO-IR -- with no header, framing, or pickle wrapper,
332 so the files are directly consumable by external NVIDIA tools
333 (``cuobjdump``, ``nvdisasm``, ``cuda-gdb``).
335 Eviction is by least-recently-*read* time: every successful read bumps
336 the entry's ``atime``, and the size enforcer evicts oldest atime
337 first.
339 .. note:: **Best-effort writes.**
341 On Windows, ``os.replace`` raises ``PermissionError`` (winerror
342 32 / 33) when another process holds the target file open. This
343 backend retries with bounded backoff (~185 ms) and, if still
344 failing, drops the cache write silently and returns success-shaped
345 control flow. The next call will see no entry and recompile. POSIX
346 and other ``PermissionError`` codes propagate.
348 .. note:: **Atomic for readers, not crash-durable.**
350 Each entry's temp file is ``fsync``-ed before ``os.replace``, but
351 the containing directory is **not** ``fsync``-ed. A host crash
352 between write and the next directory commit may lose recently
353 added entries; surviving entries remain consistent.
355 .. note:: **Cross-version sharing.**
357 The cache is safe to share across ``cuda.core`` patch releases:
358 every key produced by :func:`make_program_cache_key` encodes the
359 relevant backend/compiler/runtime fingerprints for its
360 compilation path (NVRTC entries pin the NVRTC version, NVVM
361 entries pin the libNVVM library and IR versions, PTX/linker
362 entries pin the chosen linker backend and its version -- and,
363 when the cuLink/driver backend is selected, the driver version
364 too; nvJitLink-backed PTX entries are deliberately
365 driver-version independent). Bumping ``_KEY_SCHEMA_VERSION``
366 (mixed into the digest by ``make_program_cache_key``) produces
367 new keys that don't collide with old entries: post-bump
368 lookups miss the old on-disk paths, and the orphaned files
369 are reaped on the next size-cap eviction pass. Entries are
370 stored verbatim as the compiled binary, so cross-patch sharing
371 only requires that the compiler-pinning surface above stays
372 stable -- there is no Python-pickle compatibility involved.
374 Parameters
375 ----------
376 path:
377 Directory that owns the cache. Created if missing. If omitted,
378 the OS-conventional user cache directory is used:
379 ``$XDG_CACHE_HOME/cuda-python/program-cache`` (Linux, defaulting
380 to ``~/.cache/cuda-python/program-cache``) or
381 ``%LOCALAPPDATA%\\cuda-python\\program-cache`` (Windows).
382 max_size_bytes:
383 Optional soft cap on total on-disk size. Enforced opportunistically
384 on writes; concurrent writers may briefly exceed it. Eviction is by
385 least-recently-read time (oldest ``st_atime`` first).
386 """
388 def __init__(
389 self,
390 path: str | os.PathLike[str] | None = None,
391 *,
392 max_size_bytes: int | None = None,
393 ) -> None:
394 if max_size_bytes is not None and max_size_bytes <= 0: 1rvwxybRopmqK6fezAXn#ZHITUJBWC$*+lDEakc!PigYQudF8Nhj7S1V90LO%G'M423st
395 raise ValueError("max_size_bytes must be positive or None (0 would evict every write)") 1*+
396 self._root = Path(path) if path is not None else _default_cache_dir() 1rvwxybRopmqK6fezAXn#ZHITUJBWC$lDEakc!PigYQudF8Nhj7S1V90LO%G'M423st
397 self._entries = self._root / _ENTRIES_SUBDIR 1rvwxybRopmqK6fezAXn#ZHITUJBWC$lDEakc!PigYQudF8Nhj7S1V90LO%G'M423st
398 self._tmp = self._root / _TMP_SUBDIR 1rvwxybRopmqK6fezAXn#ZHITUJBWC$lDEakc!PigYQudF8Nhj7S1V90LO%G'M423st
399 self._max_size_bytes = max_size_bytes 1rvwxybRopmqK6fezAXn#ZHITUJBWC$lDEakc!PigYQudF8Nhj7S1V90LO%G'M423st
400 # Permissions (see PR #2399):
401 # root/ and entries/ use default permissions so a shared cache (e.g. one
402 # a group shares on a cluster) keeps working. The cached files themselves
403 # are still private: each is written to tmp/ as owner-only and moved into
404 # entries/, which keeps its permissions. tmp/ is made owner-only so no one
405 # can read or swap a file while it's being written. We don't chmod, so an
406 # existing directory is left as-is.
407 # Trade-off: if a group deliberately shares a writable entries/, a member
408 # could replace a cached file. Blocking that needs a check at load time,
409 # not just permissions, and is out of scope here.
410 self._root.mkdir(parents=True, exist_ok=True) 1rvwxybRopmqK6fezAXn#ZHITUJBWC$lDEakc!PigYQudF8Nhj7S1V90LO%G'M423st
411 self._entries.mkdir(exist_ok=True) 1rvwxybRopmqK6fezAXn#ZHITUJBWC$lDEakc!PigYQudF8Nhj7S1V90LO%G'M423st
412 self._tmp.mkdir(exist_ok=True, mode=0o700) 1rvwxybRopmqK6fezAXn#ZHITUJBWC$lDEakc!PigYQudF8Nhj7S1V90LO%G'M423st
413 # Opportunistic startup sweep of orphaned temp files left by any
414 # crashed writers. Age-based so concurrent in-flight writes from
415 # other processes are preserved.
416 self._sweep_stale_tmp_files() 1rvwxybRopmqK6fezAXn#ZHITUJBWC$lDEakc!PigYQudF8Nhj7S1V90LO%G'M423st
417 # Incremental size tracker. Without it every ``__setitem__`` would
418 # walk ``entries/`` + ``tmp/`` to compute the total -- O(n) per
419 # write. With it: writes update the tracker by the net delta in O(1)
420 # and only walk on eviction (which already needs the scan to sort
421 # entries by atime). The tracker is seeded by one full scan at open
422 # time and refreshed on every eviction pass; cross-process drift
423 # (other writers/deleters) self-corrects the next time eviction
424 # fires. The lock guards mutations so multi-threaded writers in
425 # the same process don't interleave the read-modify-write on the
426 # int. Skipped entirely when ``max_size_bytes is None`` -- without
427 # a cap the tracker is dead weight.
428 self._size_lock = threading.Lock() 1rvwxybRopmqK6fezAXn#ZHITUJBWC$lDEakc!PigYQudF8Nhj7S1V90LO%G'M423st
429 self._tracked_size_bytes = self._compute_total_size() if max_size_bytes is not None else 0 1rvwxybRopmqK6fezAXn#ZHITUJBWC$lDEakc!PigYQudF8Nhj7S1V90LO%G'M423st
431 # -- key-to-path helpers -------------------------------------------------
433 def _path_for_key(self, key: object) -> Path:
434 k = _as_key_bytes(key) 1rvwxybRopmqK6fezAXnZHITUJBWClDEakcPigYQudF8Nhj7S1V0LOGM423st
435 # Hash the key to a fixed-length identifier so arbitrary-length user
436 # keys never exceed per-component filename limits (typically 255 on
437 # ext4 / NTFS).
438 #
439 # FIPS: must use a FIPS-approved hash algorithm. FIPS-enforcing
440 # systems can disable non-approved hashlib algorithms (for example
441 # blake2b) at the OpenSSL level. See #2043.
442 #
443 # With a 256-bit SHA-256 digest, the cache relies on collision
444 # resistance for key uniqueness -- two distinct keys hashing to the
445 # same path is astronomically unlikely (~2^128 practical collision
446 # work).
447 digest = hashlib.sha256(k, usedforsecurity=False).hexdigest() 1rvwxybRopmqK6fezAXnZHITUJBWClDEakcPigYQudF8Nhj7S1V0LOGM423st
448 return self._entries / digest[:2] / digest[2:] 1rvwxybRopmqK6fezAXnZHITUJBWClDEakcPigYQudF8Nhj7S1V0LOGM423st
450 # -- mapping API ---------------------------------------------------------
452 def __getitem__(self, key: object) -> bytes:
453 path = self._path_for_key(key) 1rvwxybRK6zAnHIJBClDEacuFLOGM23st
454 try: 1rvwxybRK6zAnHIJBClDEacuFLOGM23st
455 # The helper retries on Windows transient sharing-violation
456 # PermissionErrors so a racing rewriter doesn't turn a hit
457 # into a spurious propagated error.
458 st, data = _stat_and_read_with_sharing_retry(path) 1rvwxybRK6zAnHIJBClDEacuFLOGM23st
459 except FileNotFoundError: 1rbRK6HIJacust
460 raise KeyError(key) from None 1rbRK6HIJacust
461 # Bump atime to "now" so eviction (which sorts by st_atime) treats
462 # this read as the entry's most recent use. Best-effort: filesystems
463 # mounted ``noatime`` or with restrictive ACLs may refuse, in which
464 # case the cap still bounds size but eviction degrades toward FIFO
465 # rather than true LRU.
466 _touch_atime(path, st) 1rvwxybzAnBClDEacFLOGM23st
467 return data 1rvwxybzAnBClDEacFLOGM23st
469 def __setitem__(self, key: object, value: bytes | bytearray | memoryview | ObjectCode) -> None:
470 data = _extract_bytes(value) 1rvwxybRopmqKfezAXn#ZHITUJBWC$lDEakcPigYQudF8NhjS1V0LOG4st
471 target = self._path_for_key(key) 1rvwxybRopmqKfezAXnZHITUJBWClDEakcPigYQudF8NhjS1V0LOG4st
472 target.parent.mkdir(parents=True, exist_ok=True) 1rvwxybRopmqKfezAXnZHITUJBWClDEakcPigYQudF8NhjS1V0LOG4st
473 # Re-create ``tmp/`` if something deleted it after ``__init__``
474 # (operators clearing the cache by hand, ``rm -rf cache_dir/tmp``,
475 # another process's overzealous wipe). Cheap and idempotent;
476 # without it, every subsequent write would crash with
477 # FileNotFoundError even though we could trivially recover.
478 self._tmp.mkdir(parents=True, exist_ok=True) 1rvwxybRopmqKfezAXnZHITUJBWClDEakcPigYQudF8NhjS1V0LOG4st
480 # Stat the existing entry (if any) BEFORE the replace so we can
481 # update the tracker by the net delta. A racing writer that lands
482 # an ``os.replace`` between this stat and our own makes ``old_size``
483 # slightly off; the next ``_enforce_size_cap`` reconciles by
484 # re-scanning. Skipped when ``max_size_bytes is None`` (no tracker).
485 old_size = 0 1rvwxybRopmqKfezAXnZHITUJBWClDEakcPigYQudF8NhjS1V0LOG4st
486 if self._max_size_bytes is not None: 1rvwxybRopmqKfezAXnZHITUJBWClDEakcPigYQudF8NhjS1V0LOG4st
487 try: 1bfeakcigdNhj
488 old_size = target.stat().st_size 1bfeakcigdNhj
489 except FileNotFoundError: 1bfeakcigdNhj
490 old_size = 0 1bfeakcigdNhj
492 fd, tmp_name = tempfile.mkstemp(prefix="entry-", dir=self._tmp) 1rvwxybRopmqKfezAXnZHITUJBWClDEakcPigYQudF8NhjS1V0LOG4st
493 tmp_path = Path(tmp_name) 1rvwxybRopmqKfezAXnZHITUJBWClDEakcPigYQudFNhjS1V0LOG4st
494 try: 1rvwxybRopmqKfezAXnZHITUJBWClDEakcPigYQudFNhjS1V0LOG4st
495 with os.fdopen(fd, "wb") as fh: 1rvwxybRopmqKfezAXnZHITUJBWClDEakcPigYQudFNhjS1V0LOG4st
496 fh.write(data) 1rvwxybRopmqKfezAXnZHITUJBWClDEakcPigYQudFNhjS1V0LOG4st
497 fh.flush() 1rvwxybRopmqKfezAXnZHITUJBWClDEakcPigYQudFNhjS1V0LOG4st
498 os.fsync(fh.fileno()) 1rvwxybRopmqKfezAXnZHITUJBWClDEakcPigYQudFNhjS1V0LOG4st
499 # Retry os.replace under Windows sharing/lock violations; only
500 # give up (and drop the cache write) after a bounded backoff, so
501 # transient contention is not turned into a silent miss.
502 # Non-sharing PermissionErrors and all POSIX PermissionErrors
503 # propagate immediately (real config problem).
504 if not _replace_with_sharing_retry(tmp_path, target): 1rvwxybRopmqKfezAXnZHITUJBWClDEakcPigYQudFNhjS1V0LOG4st
505 with contextlib.suppress(FileNotFoundError): 1HIJ
506 tmp_path.unlink() 1HIJ
507 return 1HIJ
508 except BaseException: 1RZTU
509 with contextlib.suppress(FileNotFoundError): 1RZTU
510 tmp_path.unlink() 1RZTU
511 raise 1RZTU
513 if self._max_size_bytes is None: 1rvwxybopmqKfezAXnBWClDEakcPigYQudFNhjS1V0LOG4st
514 return 1rvwxyopmqKzAXnBWClDEPYQuFS1V0LOG4st
516 # O(1) tracker update. Only run the scan-heavy ``_enforce_size_cap``
517 # when this write actually pushes the running total above the cap.
518 new_size = len(data) 1bfeakcigdNhj
519 with self._size_lock: 1bfeakcigdNhj
520 self._tracked_size_bytes += new_size - old_size 1bfeakcigdNhj
521 over_cap = self._tracked_size_bytes > self._max_size_bytes 1bfeakcigdNhj
522 if over_cap: 1bfeakcigdNhj
523 self._enforce_size_cap() 1bfeakcgdh
525 def __delitem__(self, key: object) -> None:
526 path = self._path_for_key(key) 1KiQu7
527 # Stat before unlink so we can decrement the tracker by the actual
528 # on-disk size. Best-effort: if the file vanishes between stat and
529 # unlink (concurrent eviction), we treat the delete as a miss --
530 # matching the behaviour callers expect (KeyError) and leaving the
531 # tracker untouched (the racing eviction already accounted for it).
532 size = 0 1KiQu7
533 if self._max_size_bytes is not None: 1KiQu7
534 try: 1i7
535 size = path.stat().st_size 1i7
536 except FileNotFoundError: 17
537 raise KeyError(key) from None 17
538 try: 1KiQu
539 _unlink_with_sharing_retry(path) 1KiQu
540 except FileNotFoundError: 1KQ
541 raise KeyError(key) from None 1K
542 if self._max_size_bytes is not None: 1Kiu
543 with self._size_lock: 1i
544 # Clamp at zero. A racing ``_enforce_size_cap`` can re-seed the
545 # tracker between our stat and our subtract; if its scan ran
546 # AFTER we unlinked, its reseed value didn't include ``size``,
547 # so subtracting ``size`` again here would undercount reality
548 # by ``size``. Repeated under contention, an unclamped subtract
549 # walks the tracker negative -- and once negative, the
550 # ``tracker > cap`` check that gates ``_enforce_size_cap``
551 # never fires, so eviction dies silently and there is no
552 # self-healing path (the only reseed point is the function
553 # that no longer runs). Clamping leaves us at worst
554 # undercounting (the next reseed corrects it) instead of
555 # entering the permanently-broken negative state.
556 self._tracked_size_bytes = max(0, self._tracked_size_bytes - size) 1i
558 def __len__(self) -> int:
559 """Return the number of files currently in ``entries/``.
561 This is a count of on-disk files, not of keys reachable through
562 ``make_program_cache_key``. After a ``_KEY_SCHEMA_VERSION`` bump
563 old entries become unreachable by lookup but remain on disk
564 until eviction reaps them; ``__len__`` keeps counting them
565 until then. The same is true for entries written by callers
566 using arbitrary user keys -- the backend has no way to tell a
567 live entry from an orphan without knowing the caller's keying
568 scheme.
569 """
570 # ``_iter_entry_paths`` already filters with ``entry.is_file()``,
571 # so don't stat each path a second time here.
572 return sum(1 for _ in self._iter_entry_paths()) 1q6XnYjS1V
574 def clear(self) -> None:
575 # Snapshot stat alongside path so we can refuse to unlink an entry
576 # that was concurrently replaced by another process between the
577 # snapshot scan and the unlink. Same stat-guard contract as
578 # ``_prune_if_stat_unchanged`` and ``_enforce_size_cap``.
579 snapshot = [] 1opmqj
580 for path in self._iter_entry_paths(): 1opmqj
581 try: 1opmqj
582 snapshot.append((path, path.stat())) 1opmqj
583 except FileNotFoundError:
584 continue
585 for path, st_before in snapshot: 1opmqj
586 _prune_if_stat_unchanged(path, st_before) 1opmqj
587 # Sweep ONLY stale temp files. Deleting a young temp would race with
588 # another process between ``mkstemp`` and ``os.replace`` and turn its
589 # write into ``FileNotFoundError`` instead of a successful commit.
590 self._sweep_stale_tmp_files() 1opmqj
591 # Remove empty subdirs (best-effort; concurrent writers may re-create).
592 if self._entries.exists(): 1opmqj
593 for sub in sorted(self._entries.iterdir(), reverse=True): 1opmqj
594 if sub.is_dir(): 1opmqj
595 with contextlib.suppress(OSError): 1opmqj
596 sub.rmdir() 1opmqj
597 # The directory is now (almost) empty -- but a concurrent writer may
598 # have landed a fresh entry between the snapshot and the unlink, and
599 # young temp files were intentionally preserved. Re-derive the
600 # tracker from the post-clear state instead of zeroing blindly.
601 if self._max_size_bytes is not None: 1opmqj
602 actual = self._compute_total_size() 1j
603 with self._size_lock: 1j
604 self._tracked_size_bytes = actual 1j
606 # -- internals -----------------------------------------------------------
608 def _iter_entry_paths(self) -> Iterable[Path]:
609 # ``os.scandir`` returns ``DirEntry`` objects whose ``is_dir`` /
610 # ``is_file`` methods consult the cached dirent type from the
611 # ``readdir`` result on filesystems that report it (ext4, NTFS, ...),
612 # avoiding a per-entry ``stat`` syscall. ``Path.iterdir`` also wraps
613 # ``scandir`` but discards the cached type, forcing a separate
614 # ``stat`` for every ``Path.is_dir`` / ``Path.is_file``. The ``with``
615 # blocks release the underlying directory handle deterministically
616 # when the consumer stops early -- otherwise a leaked handle blocks
617 # deletes/renames on Windows until GC.
618 try: 1bopmq6feXnakcigYdNhj7S1VM
619 with os.scandir(self._entries) as outer: 1bopmq6feXnakcigYdNhj7S1VM
620 for sub in outer: 1bopmq6feXnakcigYdNhj7SVM
621 if not sub.is_dir(follow_symlinks=False): 1bopmqfeXnakcigYdhjSVM
622 continue 1V
623 try: 1bopmqfeXnakcigYdhjSVM
624 with os.scandir(sub.path) as inner: 1bopmqfeXnakcigYdhjSVM
625 yield from (Path(entry.path) for entry in inner if entry.is_file(follow_symlinks=False)) 1bopmqfeXnakcigYdhjSVM
626 except FileNotFoundError:
627 continue
628 except FileNotFoundError: 11
629 return 11
631 def _compute_total_size(self) -> int:
632 """Walk ``entries/`` + ``tmp/`` and return the on-disk byte total.
634 Used to seed the tracker at open time and to refresh it after every
635 eviction pass. Best-effort: files that vanish under us during the
636 walk (concurrent eviction by this or another process) are skipped.
637 Tracked total may briefly differ from this scan's result under
638 cross-process contention; the next eviction will reconcile.
639 """
640 total = 0 1bfeakcigdNhj7M
641 for path in self._iter_entry_paths(): 1bfeakcigdNhj7M
642 try: 1bakcM
643 total += path.stat().st_size 1bakcM
644 except FileNotFoundError:
645 continue
646 return total + self._sum_tmp_sizes() 1bfeakcigdNhj7M
648 def _iter_tmp_entries(self) -> Iterable[os.DirEntry[str]]:
649 # Mirror ``_iter_entry_paths``: scandir + cached d_type for the
650 # file/dir filter + deterministic handle close on early exit.
651 # Yields ``DirEntry`` (not Path) so callers can use ``entry.stat``
652 # / ``entry.path`` directly without an extra wrap.
653 try: 1rvwxybRopmqK6fezAXn#ZHITUJBWC$lDEakc!PigYQudF8Nhj7S1V90LO%G'M423st
654 with os.scandir(self._tmp) as it: 1rvwxybRopmqK6fezAXn#ZHITUJBWC$lDEakc!PigYQudF8Nhj7S1V90LO%G'M423st
655 yield from (entry for entry in it if entry.is_file(follow_symlinks=False)) 1rvwxybRopmqK6fezAXn#ZHITUJBWC$lDEakc!PigYQudF8Nhj7S1V90LO%G'M423st
656 except FileNotFoundError: 19
657 return 19
659 def _sum_tmp_sizes(self) -> int:
660 """Sum sizes of every file in ``tmp/``, skipping vanished entries.
662 Both ``_compute_total_size`` (open-time seed) and
663 ``_enforce_size_cap`` (eviction reconciliation) need this --
664 temp files occupy disk too, so undercounting them would let
665 bursts of in-flight writes silently exceed ``max_size_bytes``.
666 """
667 total = 0 1bfeakcigdNhj79M
668 for entry in self._iter_tmp_entries(): 1bfeakcigdNhj79M
669 try: 1a
670 total += entry.stat(follow_symlinks=False).st_size 1a
671 except FileNotFoundError:
672 continue
673 return total 1bfeakcigdNhj79M
675 def _sweep_stale_tmp_files(self) -> None:
676 """Remove temp files left behind by crashed writers.
678 Age threshold is conservative (``_TMP_STALE_AGE_SECONDS``) so an
679 in-flight write from another process is not interrupted. Best
680 effort: a missing file or a permission failure is ignored.
681 """
682 cutoff = time.time() - _TMP_STALE_AGE_SECONDS 1rvwxybRopmqK6fezAXn#ZHITUJBWC$lDEakc!PigYQudF8Nhj7S1V90LO%G'M423st
683 for entry in self._iter_tmp_entries(): 1rvwxybRopmqK6fezAXn#ZHITUJBWC$lDEakc!PigYQudF8Nhj7S1V90LO%G'M423st
684 try: 1oma!
685 if entry.stat(follow_symlinks=False).st_mtime < cutoff: 1oma!
686 os.unlink(entry.path) 1m!
687 except (FileNotFoundError, PermissionError):
688 continue
690 def _enforce_size_cap(self) -> None:
691 if self._max_size_bytes is None: 1bfeakcigdhS
692 return 1S
693 # Sweep stale temp files first so a long-dead writer's leftovers
694 # don't drag the apparent size up and force needless eviction.
695 self._sweep_stale_tmp_files() 1bfeakcigdh
696 entries = [] 1bfeakcigdh
697 total = 0 1bfeakcigdh
698 # Count both committed entries AND surviving temp files: temp files
699 # occupy disk too, even if they're young. Without this the soft cap
700 # silently undercounts in-flight writes.
701 #
702 # Trade-off under burst concurrency: many young temp files (each
703 # below the stale-sweep threshold) can push ``total`` above
704 # ``max_size_bytes`` with only committed entries left to evict.
705 # That can over-evict committed entries during the burst; once
706 # the burst subsides and the temps land via ``os.replace`` (or
707 # are reaped by a later sweep), the cap re-stabilises. This is
708 # consistent with the documented soft-cap contract -- callers
709 # that need a hard bound should leave the cap None and prune
710 # externally.
711 for path in self._iter_entry_paths(): 1bfeakcigdh
712 try: 1bfeakcigdh
713 st = path.stat() 1bfeakcigdh
714 except FileNotFoundError:
715 continue
716 # Carry the full stat so eviction can guard against a concurrent
717 # os.replace that swapped a fresh entry into this path between
718 # snapshot and unlink. Eviction below sorts by ``st_atime`` so
719 # entries that callers actually read recently survive
720 # write-only churn (true LRU instead of FIFO).
721 entries.append((st.st_atime, st.st_size, path, st)) 1bfeakcigdh
722 total += st.st_size 1bfeakcigdh
723 total += self._sum_tmp_sizes() 1bfeakcigdh
724 if total <= self._max_size_bytes: 1bfeakcigdh
725 # Re-seed the tracker from the scan: catches drift from
726 # cross-process writers/deleters that the per-write delta
727 # accounting wouldn't have observed. Reaching here means the
728 # tracker was over-cap but the disk truth is under-cap, so
729 # this assignment is the cheapest reconciliation point we get.
730 with self._size_lock: 1ki
731 self._tracked_size_bytes = total 1ki
732 return 1ki
733 entries.sort(key=lambda e: e[0]) # oldest atime first 1bfeacgdh
734 for _atime, size, path, st_before in entries: 1bfeacgdh
735 if total <= self._max_size_bytes: 1bfeacgdh
736 break 1bacgh
737 # _prune_if_stat_unchanged refuses if a writer replaced the file
738 # between snapshot and now, so eviction can't silently delete a
739 # freshly-committed entry from another process.
740 try: 1bfeacgdh
741 stat_now = path.stat() 1bfeacgdh
742 except FileNotFoundError:
743 total -= size
744 continue
745 if _stat_key(stat_now) != _stat_key(st_before): 1bfeacgdh
746 # File was replaced -- don't unlink, but update ``total`` to
747 # reflect the replacement's actual size or the cap check
748 # below could declare us done while still over the limit.
749 total += stat_now.st_size - size
750 continue
751 # Tolerate Windows sharing violations during eviction: another
752 # process may briefly hold the file open for a read. Skip this
753 # entry; a later eviction pass will retry. Same outcome as if
754 # the stat-guard above had triggered. Other PermissionErrors
755 # (POSIX ACL, Windows non-sharing winerrors) are real config
756 # problems -- surface them rather than silently exceed the cap.
757 try: 1bfeacgdh
758 _unlink_with_sharing_retry(path) 1bfeacgdh
759 total -= size 1bacgh
760 except FileNotFoundError: 1fed
761 pass
762 except PermissionError as exc: 1fed
763 if not _is_windows_sharing_violation(exc): 1fed
764 raise 1fe
765 # Reconcile: after the eviction pass, ``total`` reflects what we
766 # believe the disk now holds. Re-seed the tracker so the next write
767 # accumulates from a fresh baseline.
768 with self._size_lock: 1bacgdh
769 self._tracked_size_bytes = total 1bacgdh