diff options
| author | godosa <godosa@godosa.eu> | 2026-10-07 00:14:38 +0200 |
|---|---|---|
| committer | godosa <godosa@godosa.eu> | 2026-10-07 00:14:38 +0200 |
| commit | 3443c1c65e9f1753e1e656b35d08416c1fa298f2 (patch) | |
| tree | 4e43236f460145a4d75d1b4616dcb7aa6ef08f51 /servecache.py | |
| download | worldmap-viewer-3443c1c65e9f1753e1e656b35d08416c1fa298f2.tar.gz worldmap-viewer-3443c1c65e9f1753e1e656b35d08416c1fa298f2.zip | |
worldmap-viewer: initial public history
Diffstat (limited to 'servecache.py')
| -rw-r--r-- | servecache.py | 167 |
1 files changed, 167 insertions, 0 deletions
diff --git a/servecache.py b/servecache.py new file mode 100644 index 0000000..391d603 --- /dev/null +++ b/servecache.py @@ -0,0 +1,167 @@ +# map/servecache.py +"""Serve cache: the arrays the viewer server needs, written once as .npy files and memory-mapped afterwards, so the +world and every refined area stay on disk instead of in RAM.""" +from __future__ import annotations + +import json +import os +import shutil +import time +from pathlib import Path + +import numpy as np + +import dedupe + +VERSION = 1 # bump when the layout of a cache folder changes +PREFIX = "serve-" +DONE = "done" + + +def low_memory() -> bool: + """WORLDGEN_LOW_MEMORY=1 (or --low-memory): build caches with fields on disk; slower, same values.""" + return os.environ.get("WORLDGEN_LOW_MEMORY", "") not in ("", "0") + + +def more_files() -> None: + """Each memory-mapped array holds an open file (mmap keeps one): raise the soft open-file limit to the hard one, + so a cache with hundreds of arrays loads under a low default (systemd units: 1024).""" + try: + import resource + soft, hard = resource.getrlimit(resource.RLIMIT_NOFILE) + want = hard if hard != resource.RLIM_INFINITY else max(soft, 1 << 20) + if soft != resource.RLIM_INFINITY and soft < want: + resource.setrlimit(resource.RLIMIT_NOFILE, (want, hard)) + except (ImportError, ValueError, OSError): + pass + + +class _Stream: + """A file seen only as something with write(): numpy then writes the array in chunks instead of tofile(), which + preallocates the file (fallocate) — and btrfs never compresses preallocated space. Same bytes as np.save.""" + def __init__(self, fh): + self.write = fh.write + + +def save_npy(path, v) -> None: + """np.save(path, v), byte for byte, written so a compressing file system (btrfs compress=zstd) can compress it: + the cache's mostly smooth or constant fields take about a third of the disk; memory maps read it as before.""" + with open(path, "wb") as fh: + np.lib.format.write_array(_Stream(fh), np.ascontiguousarray(v), allow_pickle=False) + + +def spiller(spill): + """keep(name, array): the array itself, or (spill folder given) saved there and handed back as a read-only memory + map, so building a big cache needs about one field of RAM at a time instead of all of them.""" + if spill is None: + return lambda name, v: v + spill = Path(spill) + more_files() + + def keep(name, v): + f = spill / f"{name}.npy" + save_npy(f, v) + return np.load(f, mmap_mode="r", allow_pickle=False) + return keep + + +def cache_dir(parent: Path, fingerprint: str) -> Path: + return Path(parent) / f"{PREFIX}{fingerprint}-v{VERSION}" + + +def write(d: Path, arrays: dict, meta: dict) -> None: + """Write arrays (one .npy each) and meta (JSON) atomically: a temporary folder, `done` last, renamed into place. + If another process finished the same folder first, theirs is kept.""" + d = Path(d) + tmp = d.parent / f"{d.name}.tmp-{os.getpid()}" + shutil.rmtree(tmp, ignore_errors=True) + tmp.mkdir(parents=True) + try: + for k, v in arrays.items(): + save_npy(tmp / f"{k}.npy", v) + (tmp / "meta.json").write_text(json.dumps({**meta, "keys": sorted(arrays)})) + (tmp / DONE).write_text("ok") + if d.exists(): + shutil.rmtree(tmp, ignore_errors=True) + return + try: + os.rename(tmp, d) + except OSError: + if read(d) is None: + raise + shutil.rmtree(tmp, ignore_errors=True) # another server finished it first + return + except BaseException: + shutil.rmtree(tmp, ignore_errors=True) + raise + dedupe.dedupe_dirs(d, peers(d)) # other eras' caches: share their identical blocks + + +def peers(d: Path) -> list: + """Complete caches of the other eras of the same regions (<regions>/<era key>/serve-*): mostly the same bytes.""" + d = Path(d) + return [p for p in sorted(d.parent.parent.glob(f"*/{PREFIX}*")) if p.parent != d.parent and (p / DONE).exists()] + + +def read(d: Path): + """(arrays as read-only memory maps, meta) of a complete cache folder, else None.""" + d = Path(d) + if not (d / DONE).exists(): + return None + more_files() + try: + meta = json.loads((d / "meta.json").read_text()) + arrays = {k: np.load(d / f"{k}.npy", mmap_mode="r", allow_pickle=False) for k in meta["keys"]} + except Exception: # noqa: BLE001 — damaged (truncated, empty): as none + return None + return arrays, meta + + +def load_or_build(d: Path, build, log=print): + """The cached arrays, or build() → (arrays, meta) written for next time and handed back as memory maps. + If the folder can't be written (read-only or full disk), the built arrays are used from RAM, with one warning.""" + got = read(d) + if got is not None: + return got + if Path(d).exists(): # complete but unreadable (e.g. power lost after the + shutil.rmtree(d, ignore_errors=True) # rename): replaced, not kept forever + arrays, meta = build() + try: + write(d, arrays, meta) + except OSError as e: + log(f"serve cache: can't write {d} ({e}); keeping the data in memory") + return arrays, meta + got = read(d) + return got if got is not None else (arrays, meta) + + +def _alive(pid: str) -> bool: + try: + os.kill(int(pid), 0) + except ProcessLookupError: + return False + except (PermissionError, ValueError, OverflowError): + return True + return True + + +def prune(parent: Path, keep, min_tmp_age_s: float = 3600.0) -> list: + """Remove serve cache folders in parent whose names aren't in keep, and unfinished temporary folders older than + min_tmp_age_s or whose writer is gone (a younger one may belong to a server still writing it). Other folders are + never touched.""" + parent, names = Path(parent), {Path(k).name for k in keep} + gone = [] + if not parent.is_dir(): + return gone + now = time.time() + for p in sorted(parent.iterdir()): + if not p.is_dir() or not p.name.startswith(PREFIX) or p.name in names: + continue + try: + if ".tmp-" in p.name and _alive(p.name.rsplit(".tmp-", 1)[1]) and now - p.stat().st_mtime < min_tmp_age_s: + continue + except FileNotFoundError: # another server removed it meanwhile + continue + shutil.rmtree(p, ignore_errors=True) + gone.append(p) + return gone |
