raw ยท 6532 bytes
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 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 |