Coverage for gws-app/gws/gis/cache/seed.py: 94%
137 statements
« prev ^ index » next coverage.py v7.16.2, created at 2026-10-05 13:35 +0200
« prev ^ index » next coverage.py v7.16.2, created at 2026-10-05 13:35 +0200
1"""Cache seeding."""
3import threading
4from typing import cast
6import gws
7import gws.lib.grid
8import gws.lib.osx
10from . import core
12DEFAULT_BLOCK_SIZE = 8
13PROGRESS_INTERVAL = 5
16def seed(root: gws.Root, opts: core.SeedOptions) -> core.SeedResult:
17 """Fill the selected caches with missing tiles.
19 Tiles are requested from the grabbers in blocks by ``opts.concurrency`` worker threads,
20 round-robin over the caches, until all blocks are done, the time limit is reached
21 or the run is interrupted. Only one seeding run can be active on the server;
22 if another one is running, nothing is done and the status is ``locked``.
24 Args:
25 root: Configuration root.
26 opts: Seeding options. Missing values default to 600 seconds and one thread.
28 Returns:
29 The seeding result with per-level statistics.
30 """
32 defaults = core.SeedOptions(filter=None, maxTime=600, concurrency=1)
33 opts = cast(core.SeedOptions, gws.u.merge(defaults, opts))
34 try:
35 with gws.u.server_lock('seed', 0):
36 return _run(root, opts)
37 except gws.LockBusyError:
38 gws.log.info('seed: already running')
39 return core.SeedResult(caches=[], seedTime=0, seedStatus='locked')
42##
45def _run(root: gws.Root, opts: core.SeedOptions) -> core.SeedResult:
46 inv = core.inventory(root)
47 if opts.filter:
48 core.apply_filter(inv, opts.filter)
49 if opts.maxAge is not None:
50 for c in inv.caches:
51 c.grabber.store.maxAge = min(opts.maxAge, c.grabber.cache.maxAge)
52 core.add_stats(inv)
53 res = core.SeedResult(caches=inv.caches, seedTime=0, seedStatus='')
54 ts = gws.u.stime()
56 queue = _BlockQueue(res.caches)
58 deadline = gws.u.stime() + opts.maxTime
59 threads = [threading.Thread(target=_worker, args=(queue, deadline), daemon=True) for _ in range(opts.concurrency)]
61 for t in threads:
62 t.start()
63 try:
64 for t in threads:
65 while t.is_alive():
66 t.join(0.5)
67 except KeyboardInterrupt:
68 gws.log.info('seed: interrupted')
69 queue.stop('interrupted')
71 queue.report()
73 res.seedTime = gws.u.stime() - ts
74 res.seedStatus = queue.seedStatus
76 return res
79class _BlockGenerator:
80 """Generates the blocks of one cache, level by level."""
82 def __init__(self, cache: core.Cache, levels: list[core.Level]):
83 self.cache = cache
84 self.levels = {lv.z: lv for lv in levels}
85 self.size = getattr(self.cache.grabber, 'requestTiles', 0) or DEFAULT_BLOCK_SIZE
86 self.blocks = self.iter_blocks()
87 self.startTime: dict[int, float] = {}
89 def iter_blocks(self):
90 """Yield tile ranges of block size, aligned to multiples of the block size."""
91 for z in sorted(self.levels):
92 x0, y0, x1, y1, _ = self.levels[z].gridRange
93 n = self.size
94 for by in range((y0 // n) * n, y1 + 1, n):
95 for bx in range((x0 // n) * n, x1 + 1, n):
96 yield max(bx, x0), max(by, y0), min(bx + n - 1, x1), min(by + n - 1, y1), z
99class _BlockQueue:
100 """Yields blocks to seed, round-robin over caches, so that no single source gets all threads."""
102 def __init__(self, caches: list[core.Cache]):
103 self.lock = threading.Lock()
104 self.caches = caches
105 self.generators: list[_BlockGenerator] = []
106 self.seedStatus = ''
107 self.stopped = False
108 self.lastReport = gws.u.stime()
110 for c in caches:
111 if c.levels:
112 self.generators.append(_BlockGenerator(c, c.levels))
114 def next_block(self) -> tuple[_BlockGenerator, gws.MapTileRange] | None:
115 """Return the next block and its generator, or ``None`` if all blocks are done."""
116 with self.lock:
117 while self.generators:
118 bg = self.generators.pop(0)
119 block = next(bg.blocks, None)
120 if block is None:
121 continue
122 z = block[4]
123 if z not in bg.startTime:
124 bg.startTime[z] = gws.u.stime()
125 bg.levels[z].cachedTiles = 0
126 self.generators.append(bg)
127 return bg, block
129 def block_complete(self, bg: _BlockGenerator, block: gws.MapTileRange, present: int, fetched: int, failed: int):
130 """Add the counts of a completed block to its level and report progress periodically."""
131 with self.lock:
132 z = block[4]
133 lv = bg.levels[z]
134 lv.cachedTiles += present
135 lv.fetchedTiles += fetched
136 lv.failedTiles += failed
137 lv.seedTime = gws.u.stime() - bg.startTime[z]
138 if gws.u.stime() - self.lastReport >= PROGRESS_INTERVAL:
139 self.report()
141 def stop(self, status: str):
142 """Stop seeding and set the status of the run and of the unfinished caches."""
143 with self.lock:
144 self.seedStatus = status
145 self.stopped = True
146 for bg in self.generators:
147 bg.cache.seedStatus = status
149 def report(self):
150 """Log the cached percentages of all caches."""
151 self.lastReport = gws.u.stime()
152 percents = {c.name: core.percentage_by_level(c) for c in self.caches}
153 for name, ps in sorted(percents.items()):
154 gws.log.info(f'seed {name}: %% {" ".join(f"{z}:{p}" for z, p in enumerate(ps))}')
156def _worker(queue: _BlockQueue, deadline: float):
157 while not queue.stopped:
158 if gws.u.stime() >= deadline:
159 gws.log.info('seed: time limit reached')
160 queue.stop('timeout')
161 return
163 p = queue.next_block()
164 if not p:
165 return
166 bg, block = p
168 try:
169 present, fetched, failed = _seed_block(queue, bg, block)
170 except Exception as exc:
171 gws.log.error(f'seed {bg.cache.name}: block {block} error: {exc!r}')
172 present, fetched, failed = 0, 0, 0
174 if not queue.stopped:
175 queue.block_complete(bg, block, present, fetched, failed)
178def _seed_block(queue: _BlockQueue, bg: _BlockGenerator, block: gws.MapTileRange) -> tuple[int, int, int]:
179 present = 0
180 missing = []
182 for mt in gws.lib.grid.enum_tiles(block):
183 if bg.cache.grabber.store.has(mt):
184 present += 1
185 else:
186 missing.append(mt)
188 if not missing:
189 return present, 0, 0
191 try:
192 for mt in missing:
193 if queue.stopped:
194 return present, 0, 0
195 bg.cache.grabber.get_tile_as_bytes(mt)
196 return present, len(missing), 0
197 except Exception as exc:
198 gws.log.warning(f'seed {bg.cache.name}: block {block} failed: {exc!r}')
199 return present, 0, len(missing)