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

1"""Cache seeding.""" 

2 

3import threading 

4from typing import cast 

5 

6import gws 

7import gws.lib.grid 

8import gws.lib.osx 

9 

10from . import core 

11 

12DEFAULT_BLOCK_SIZE = 8 

13PROGRESS_INTERVAL = 5 

14 

15 

16def seed(root: gws.Root, opts: core.SeedOptions) -> core.SeedResult: 

17 """Fill the selected caches with missing tiles. 

18 

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``. 

23 

24 Args: 

25 root: Configuration root. 

26 opts: Seeding options. Missing values default to 600 seconds and one thread. 

27 

28 Returns: 

29 The seeding result with per-level statistics. 

30 """ 

31 

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') 

40 

41 

42## 

43 

44 

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() 

55 

56 queue = _BlockQueue(res.caches) 

57 

58 deadline = gws.u.stime() + opts.maxTime 

59 threads = [threading.Thread(target=_worker, args=(queue, deadline), daemon=True) for _ in range(opts.concurrency)] 

60 

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') 

70 

71 queue.report() 

72 

73 res.seedTime = gws.u.stime() - ts 

74 res.seedStatus = queue.seedStatus 

75 

76 return res 

77 

78 

79class _BlockGenerator: 

80 """Generates the blocks of one cache, level by level.""" 

81 

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] = {} 

88 

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 

97 

98 

99class _BlockQueue: 

100 """Yields blocks to seed, round-robin over caches, so that no single source gets all threads.""" 

101 

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() 

109 

110 for c in caches: 

111 if c.levels: 

112 self.generators.append(_BlockGenerator(c, c.levels)) 

113 

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 

128 

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() 

140 

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 

148 

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))}') 

155 

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 

162 

163 p = queue.next_block() 

164 if not p: 

165 return 

166 bg, block = p 

167 

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 

173 

174 if not queue.stopped: 

175 queue.block_complete(bg, block, present, fetched, failed) 

176 

177 

178def _seed_block(queue: _BlockQueue, bg: _BlockGenerator, block: gws.MapTileRange) -> tuple[int, int, int]: 

179 present = 0 

180 missing = [] 

181 

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) 

187 

188 if not missing: 

189 return present, 0, 0 

190 

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)