worldmap-viewer

git clone https://git.godosa.eu/worldmap-viewer

master

raw ยท 6532 bytes

# 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