aboutsummaryrefslogtreecommitdiffziptar.gz
path: root/servecache.py
blob: 391d60395b187bd84214137873456ee2a73a4ded (plain)
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