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
|