Skip to content

Commit 1285756

Browse files
committed
perf: tile encoding uses all CPUs, not just --workers; robustness: async SSH sync (live log tail during transfer), per-line chunk id in log, local disk guard before SSH copy, skip .vrt, tolerate files purged mid-sync
1 parent c8c5b65 commit 1285756

3 files changed

Lines changed: 374 additions & 22 deletions

File tree

lidar2map.py

Lines changed: 36 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -2032,6 +2032,14 @@ def __init__(self, log_path):
20322032
# de tuiles, workers). Sans lui, la machine à états _buf/_cr_buf
20332033
# s'entrelace entre threads → lignes de log corrompues.
20342034
self._lock = threading.Lock()
2035+
self._chunk_actuel = None # cf. definir_chunk
2036+
2037+
def definir_chunk(self, cle):
2038+
"""Marque le chunk en cours (découpage à priori) : préfixe chaque
2039+
ligne de log qui suit, pour qu'un extrait du log soit auto-suffisant
2040+
sans avoir à remonter chercher le dernier « ── Ombrage XXX ── »."""
2041+
with self._lock:
2042+
self._chunk_actuel = cle
20352043

20362044
def _log_line(self, line):
20372045
"""Écrit une ligne dans le fichier log avec horodatage."""
@@ -2041,7 +2049,8 @@ def _log_line(self, line):
20412049
line = line.strip()
20422050
if line:
20432051
ts = time.strftime("%H:%M:%S")
2044-
self._log.write(f"[{ts}] {line}\n")
2052+
_cle = f"[{self._chunk_actuel}] " if self._chunk_actuel else ""
2053+
self._log.write(f"[{ts}] {_cle}{line}\n")
20452054

20462055
def write(self, msg):
20472056
# ── Terminal ─────────────────────────────────────────────────────────
@@ -2236,6 +2245,15 @@ def _excepthook(exc_type, exc_value, exc_tb):
22362245

22372246
_activer_log()
22382247

2248+
2249+
def _definir_chunk_log(cle):
2250+
"""Signale à _TeeLogger le chunk en cours (découpage à priori), pour
2251+
préfixer chaque ligne du log qui suit (cf. _TeeLogger.definir_chunk).
2252+
No-op si le log n'est pas actif (tests, sys.stdout non remplacé)."""
2253+
_t = sys.stdout
2254+
if isinstance(_t, _TeeLogger):
2255+
_t.definir_chunk(cle)
2256+
22392257
# ── Requêtes HTTP via urllib (stdlib, zéro dépendance) ──────────────────────
22402258
_HTTP_UA = "lidar2map/1.0 (IGN WMTS/WMS)"
22412259

@@ -8453,6 +8471,17 @@ def _warped_3857_valide(chemin):
84538471
return False
84548472

84558473

8474+
def _tile_workers_defaut():
8475+
"""Parallélisme de l'encodage de tuiles (JPEG/PNG, Pillow libère le GIL,
8476+
cf. le pool dans generer_mbtiles_lidar) : DÉCOUPLÉ de --workers, qui
8477+
plafonne les téléchargements réseau (throttle IGN ~3 en simultané, cf.
8478+
--laz-parallel pour la même logique côté conversion LAZ). L'encodage est
8479+
100% CPU local, sans lien avec ce plafond réseau : le limiter au nombre
8480+
de workers de download le bridait sans raison (ex. --workers 3 sur une
8481+
VM à 16 vCPU = 13 coeurs inutilisés pendant tout le tuilage)."""
8482+
return os.cpu_count() or 4
8483+
8484+
84568485
def generer_mbtiles_lidar(tif_source, dossier_ville, nom_ville,
84578486
zoom_min=13, zoom_max=17, format_tuiles="auto",
84588487
jpeg_quality=85, bbox_natif=None, tampon_coin_max_m=0,
@@ -11323,7 +11352,7 @@ def _tuiler_tifs_ombrages(args, tifs, dossier_ville, nom_zone, bbox,
1132311352
jpeg_quality=args.qualite_image,
1132411353
bbox_natif=bbox, tampon_coin_max_m=tampon_coin_max_m,
1132511354
ecraser_tuiles=args.tuiles_ecraser,
11326-
tile_workers=args.workers)
11355+
tile_workers=_tile_workers_defaut())
1132711356
else:
1132811357
print(f" Existing MBTiles: {mbt_path.name}, direct split/conversion")
1132911358
mbt_out = mbt_path
@@ -12405,7 +12434,7 @@ def _est_du_provider(nom):
1240512434
bbox_natif=bbox,
1240612435
source_already_warped=getattr(args, "_source_already_warped", False),
1240712436
ecraser_tuiles=_ecraser_l,
12408-
tile_workers=args.workers)
12437+
tile_workers=_tile_workers_defaut())
1240912438
elif _mbt_path.exists():
1241012439
print(f" Existing MBTiles: {_mbt_path.name}, direct split/conversion")
1241112440
_mbt_out = _mbt_path
@@ -13603,6 +13632,7 @@ def _run_split_priori(args, sous_zones, mode_desc, nom_zone, racine_pr,
1360313632

1360413633
print(f"\n ── Chunk {cle} ({i_z+1}/{n_total}) {nom_z} ──")
1360513634
print(f" {entete_chunk(coords)}")
13635+
_definir_chunk_log(cle)
1360613636
manifeste.debut_morceau(cle, nom_z)
1360713637
t0_z = time.time()
1360813638
try:
@@ -14001,7 +14031,7 @@ def _traiter_bbox_lidar_tuilage(args, bbox_natif, nom_z, nom_zone_base, manifest
1400114031
jpeg_quality=args.qualite_image,
1400214032
bbox_natif=bbox, tampon_coin_max_m=TAMPON_MAX_M,
1400314033
ecraser_tuiles=args.tuiles_ecraser,
14004-
tile_workers=args.workers)
14034+
tile_workers=_tile_workers_defaut())
1400514035
else:
1400614036
print(f" Existing MBTiles: {mbt_path.name}, direct split/conversion")
1400714037
mbt_out = mbt_path
@@ -14108,6 +14138,7 @@ def _etape_ombrage(sz):
1410814138
cle, 0, n_total)
1410914139
print(f"\n ── Ombrage {cle} {nom_z} ──")
1411014140
print(f" {entete_chunk(tuple(sz[2:]))}")
14141+
_definir_chunk_log(cle)
1411114142
manifeste.debut_morceau(cle, nom_z)
1411214143
t0 = time.time()
1411314144
dalles_precharge = prefetch.recuperer(cle)
@@ -14124,6 +14155,7 @@ def _etape_tuilage(sz):
1412414155
if manifeste.deja_traite(cle_t) and not overwrite_actif:
1412514156
return
1412614157
nom_z = f"{nom_zone}_{cle}"
14158+
_definir_chunk_log(cle_t)
1412714159
manifeste.debut_morceau(cle_t, nom_z)
1412814160
t0 = time.time()
1412914161
_traiter_bbox_lidar_tuilage(args, tuple(sz[2:]), nom_z, nom_zone, manifeste, cle,

tests/_test_rlidar2map_CLI.py

Lines changed: 199 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -562,6 +562,111 @@ def kill(self):
562562
self.assertTrue(ok)
563563
self.assertEqual(destination.read_bytes(), b"complete-new")
564564

565+
def test_scp_missing_remote_file_is_skipped_not_fatal(self):
566+
# Le pipeline distant purge ses intermédiaires (VRT voisin, etc.)
567+
# pendant que le sync tourne : un fichier de l'inventaire peut avoir
568+
# disparu au moment de la copie. Vécu : FileNotFoundError distante
569+
# plantait tout le lot. Le fichier manquant doit juste être ignoré,
570+
# les AUTRES fichiers du même lot doivent quand même passer.
571+
with tempfile.TemporaryDirectory() as tmp:
572+
parsed = rlidar2map_cli.parse_options(
573+
["--local-dir", tmp, "--sync-method", "scp", "vm.example"]
574+
)
575+
controller = rlidar2map_cli.VmController(parsed)
576+
state = remote_state("succeeded", exit_code=0)
577+
local_results = Path(tmp) / state.run_id / "results"
578+
local_results.mkdir(parents=True)
579+
580+
# PAS .vrt : exclu inconditionnellement en amont (cf.
581+
# ALWAYS_EXCLUDED_EXTENSIONS) donc jamais demandé au tout, ce qui
582+
# court-circuiterait le chemin de résilience testé ici. Un autre
583+
# intermédiaire (warp temporaire) illustre le même risque de
584+
# purge concurrente pour un fichier non filtré en amont.
585+
gone_relative = "gar9_002x003/warped_3857_tmp.tif"
586+
gone_encoded = rlidar2map_cli.base64.b64encode(
587+
gone_relative.encode("utf-8")
588+
).decode("ascii")
589+
ok_relative = "gar9_002x003/result.mbtiles"
590+
ok_encoded = rlidar2map_cli.base64.b64encode(
591+
ok_relative.encode("utf-8")
592+
).decode("ascii")
593+
ok_content = b"complete-mbtiles"
594+
595+
inventory = {
596+
gone_relative: (gone_encoded, (1, 2, 999, 3)),
597+
ok_relative: (ok_encoded, (1, 2, len(ok_content), 3)),
598+
}
599+
600+
def frame(payload):
601+
raw = json.dumps(
602+
payload, separators=(",", ":")
603+
).encode("ascii")
604+
return rlidar2map_cli.struct.pack(">Q", len(raw)) + raw
605+
606+
class FakeProcess:
607+
def __init__(self, output, returncode):
608+
self.stdin = io.BytesIO()
609+
self.stdout = io.BytesIO(output)
610+
self._final_returncode = returncode
611+
self._returncode = None
612+
613+
def wait(self):
614+
if self._returncode is None:
615+
self._returncode = self._final_returncode
616+
return self._returncode
617+
618+
def poll(self):
619+
return self._returncode
620+
621+
def kill(self):
622+
self._returncode = -9
623+
624+
digest = rlidar2map_cli.hashlib.sha256(ok_content).hexdigest()
625+
# Ordre = ordre du dict inventory (Python 3.7+ garantit l'ordre
626+
# d'insertion) : "missing" pour le .vrt, puis file/trailer normal
627+
# pour le .mbtiles. Code de sortie 0 (pas 74/unstable) : rien n'a
628+
# été lu de façon instable, le fichier était juste absent.
629+
stream = (
630+
rlidar2map_cli.FILE_STREAM_MAGIC
631+
+ frame({"type": "missing", "path": gone_encoded})
632+
+ frame(
633+
{
634+
"type": "file",
635+
"path": ok_encoded,
636+
"size": len(ok_content),
637+
}
638+
)
639+
+ ok_content
640+
+ frame(
641+
{
642+
"type": "trailer",
643+
"path": ok_encoded,
644+
"stable": True,
645+
"sha256": digest,
646+
}
647+
)
648+
+ frame({"type": "end", "count": 2})
649+
)
650+
with mock.patch.object(
651+
controller,
652+
"_remote_results_inventory",
653+
return_value=inventory,
654+
), mock.patch.object(
655+
rlidar2map_cli.subprocess,
656+
"Popen",
657+
return_value=FakeProcess(stream, 0),
658+
):
659+
ok = controller._sync_results_scp(state, local_results)
660+
661+
self.assertTrue(ok)
662+
self.assertEqual(
663+
(local_results / ok_relative).read_bytes(), ok_content
664+
)
665+
self.assertFalse((local_results / gone_relative).exists())
666+
self.assertFalse(
667+
list(local_results.parent.glob(".rlidar2map-sync-*"))
668+
)
669+
565670
def test_purge_markers_are_monotonic_across_stale_monitors(self):
566671
with tempfile.TemporaryDirectory() as tmp:
567672
parsed = rlidar2map_cli.parse_options(
@@ -663,10 +768,12 @@ def fake_run(command, **kwargs):
663768
self.assertNotEqual(calls[0][0], "scp")
664769

665770
def test_sync_only_excludes_mapping(self):
771+
# .vrt exclu inconditionnellement (intermédiaire jamais utile en
772+
# local, cf. ALWAYS_EXCLUDED_EXTENSIONS), quel que soit --sync-only.
666773
for value, expected in (
667-
("tout", ()),
668-
("ombrages", (".mbtiles", ".rmap", ".sqlitedb")),
669-
("carte", (".tif",)),
774+
("tout", (".vrt",)),
775+
("ombrages", (".vrt", ".mbtiles", ".rmap", ".sqlitedb")),
776+
("carte", (".vrt", ".tif")),
670777
):
671778
parsed = rlidar2map_cli.parse_options(
672779
["--sync-only", value, "vm.example"]
@@ -741,6 +848,95 @@ def fake_run(command, **kwargs):
741848
self.assertEqual(method, "ssh")
742849
self.assertEqual(calls, [])
743850

851+
def test_scp_transfer_skipped_when_local_disk_is_low(self):
852+
# Garde-fou disque local avant une copie SSH : espace insuffisant ->
853+
# annulation propre (retry au prochain cycle), pas de tentative de
854+
# copie qui échouerait à mi-chemin ou saturerait le disque local.
855+
with tempfile.TemporaryDirectory() as tmp:
856+
parsed = rlidar2map_cli.parse_options(
857+
["--local-dir", tmp, "--sync-method", "auto", "vm.example"]
858+
)
859+
deps = rlidar2map_cli.RuntimeDeps(which=lambda _name: None)
860+
controller = rlidar2map_cli.VmController(parsed, deps)
861+
state = remote_state("succeeded", exit_code=0)
862+
calls = []
863+
864+
def fake_run(command, **kwargs):
865+
calls.append((list(command), kwargs))
866+
return subprocess.CompletedProcess(command, 0, b"", b"")
867+
868+
encoded = rlidar2map_cli.base64.b64encode(
869+
b"gar9_001x001_svf_flux.tif"
870+
).decode("ascii")
871+
inventory = {
872+
"gar9_001x001_svf_flux.tif": (encoded, (1, 2, 5_000_000_000, 4)),
873+
}
874+
fake_usage = mock.Mock(free=1_000) # bien en dessous de la marge
875+
local_results = controller.local_run_dir(state) / "results"
876+
local_results.mkdir(parents=True, exist_ok=True)
877+
with mock.patch.object(
878+
rlidar2map_cli.subprocess, "run", fake_run
879+
), mock.patch.object(
880+
controller,
881+
"_remote_results_inventory",
882+
return_value=inventory,
883+
), mock.patch.object(
884+
rlidar2map_cli.shutil, "disk_usage", return_value=fake_usage,
885+
):
886+
# _sync_results_scp directement (pas sync_once) : isole le
887+
# garde-fou disque du transfert du log final, sans rapport.
888+
ok = controller._sync_results_scp(state, local_results)
889+
890+
self.assertFalse(ok)
891+
self.assertEqual(calls, [])
892+
893+
def test_print_remote_log_tail_collapses_cr_progress_across_polls(self):
894+
# run.log est la capture brute (tmux/tee) : ses répétitions \r ne
895+
# doivent PAS devenir des lignes [VM] distinctes (splitlines() coupe
896+
# aussi sur \r) - seul l'état final d'une répétition compte, et
897+
# l'état doit survivre à la frontière entre deux sondages (tail
898+
# incrémental), une répétition pouvant être coupée en plein milieu.
899+
with tempfile.TemporaryDirectory() as tmp:
900+
parsed = rlidar2map_cli.parse_options(
901+
["--local-dir", tmp, "--sync-method", "auto", "vm.example"]
902+
)
903+
deps = rlidar2map_cli.RuntimeDeps(which=lambda _name: None)
904+
controller = rlidar2map_cli.VmController(parsed, deps)
905+
state = remote_state()
906+
907+
# Sondage 1 : 3 répétitions \r, la dernière coupée avant tout \r/\n
908+
# (poll suivant en pleine barre de progression).
909+
chunk1 = (
910+
b"SVF chunked: 10% (1/10)\r"
911+
b"SVF chunked: 20% (2/10)\r"
912+
b"SVF chunked: 30% (3/10)"
913+
)
914+
# Sondage 2 : continuation -> \r finalise le 30%, PUIS la vraie
915+
# ligne de fin, terminée par \n.
916+
chunk2 = b"\rSVF chunked: done (10 blocks)\n"
917+
918+
results = iter([chunk1, chunk2])
919+
920+
def fake_run(_command, **_kwargs):
921+
return subprocess.CompletedProcess(
922+
_command, 0, next(results), b""
923+
)
924+
925+
printed = io.StringIO()
926+
with mock.patch.object(
927+
rlidar2map_cli.subprocess, "run", fake_run
928+
), mock.patch("sys.stdout", printed):
929+
controller.print_remote_log_tail(state)
930+
controller.print_remote_log_tail(state)
931+
932+
output = printed.getvalue()
933+
self.assertNotIn("10%", output)
934+
self.assertNotIn("20%", output)
935+
self.assertNotIn("30%", output)
936+
self.assertEqual(
937+
output.count("SVF chunked: done (10 blocks)"), 1
938+
)
939+
744940

745941
class ControllerTests(unittest.TestCase):
746942
def _controller_mock(self, states, sync_results):

0 commit comments

Comments
 (0)