Files
lidar_rendu/lidar_pipeline/webapp.py

1551 lines
66 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

"""Serveur web de la carte LiDAR : sert l'index et expose l'API de génération.
Lancé via `./run.sh --serve [PORT]` (input/ monté en écriture pour permettre
le téléchargement IGN). Endpoints :
GET / → carte interactive (output/index.html)
GET /api/status → état de la génération en cours (ou dernière terminée)
GET /api/tiles → données de carte (index_tiles.json) pour l'affichage
en direct ; ?stamp=X → réponse allégée si inchangé
POST /api/sync → régénérer l'index local (vignettes + carte) en
arrière-plan (les tuiles sont servies à la demande)
POST /api/preview → cellules 1 km intersectant une bbox WGS84, OU (option
all_missing=true) toutes les dalles LHD présentes dans
input/ qui manquent au moins une visualisation demandée.
Sont incluses les tuiles absentes ET les tuiles
existantes incomplètes ; option viz = noms d'étapes
demandées, option regenerate=true pour inclure les
tuiles complètes
POST /api/generate → télécharge (géoplateforme IGN) puis traite des cellules
(option regenerate=true ajoute --force : les visualisations
sont refaites mais la classification existante est
conservée ; option reclassify=true ajoute en plus
--force-classification pour relancer la classification
du sol — sinon changer simplement ground_class
reclassifie déjà les tuiles dont la méthode diffère ;
option all_missing=true traite toutes les dalles présentes
dans input/ qui manquent les visualisations demandées, sans
téléchargement ; sans --force, le pipeline ne génère que
les visualisations manquantes des tuiles existantes)
POST /api/stop → arrête la génération en cours (SIGTERM au pipeline, qui
nettoie ses workers et processus PDAL ; SIGKILL du groupe
en repli après 15 s)
POST /api/export → assemble des dalles adjacentes en une image (PNG/JPEG/
WebP) ou un PDF multi-couches, consultable sur téléphone
(module export.py — mosaïque sans couture, habillage
titre/échelle/nord) ; GET /api/export/file/{nom} sert le
fichier en téléchargement. Marche aussi sur la webapp
légère : les tuiles viennent du cache local, aucune
délégation à la machine de traitement.
GET /api/presets → presets de couches personnalisés (couches visibles,
opacités, ordre, fond OSM) stockés sur le serveur et
partagés par tous les navigateurs connectés
POST /api/presets → créer/remplacer un preset de couches (même identifiant
= remplacement)
DELETE /api/presets/{id} → supprimer un preset de couches
Fichiers statiques : /assets (interface), /index_thumbs, /index_subtiles,
/visualisations, /DTM.
Un seul job à la fois : la génération lance `python -m lidar_pipeline` en
sous-processus avec --fetch-tiles + --file, journalisé dans .generation.log.
Architecture deux machines (cf. docs/DEPLOY_WEBAPP.md) : la webapp légère
(sans PDAL/GPU, ex. Raspberry Pi) délègue la génération à la machine de
traitement via LIDAR_GENERATION_URL — /api/generate, /api/preview et
/api/status sont transmis tels quels au service distant, qui exécute la même
webapp avec le pipeline complet. Les images produites sont servies à la
demande depuis GENERATION_URL (cache local peuplé au fil des consultations)
puis /api/sync régénère vignettes et carte localement. LIDAR_API_TOKEN
(machine de traitement) + LIDAR_REMOTE_TOKEN (webapp légère) protègent
optionnellement les appels distants.
"""
import json
import logging
import math
import os
import re
import signal
import subprocess
import sys
import threading
import time
import urllib.parse
import urllib.request
import uuid
from pathlib import Path
from typing import Optional
from fastapi import Depends, FastAPI, Header, HTTPException, Request
from fastapi.responses import FileResponse, JSONResponse
from pydantic import BaseModel, Field
INPUT_DIR = Path(os.environ.get("LIDAR_INPUT_DIR", "/data/input"))
OUTPUT_DIR = Path(os.environ.get("LIDAR_OUTPUT_DIR", "/data/output"))
PORT = int(os.environ.get("LIDAR_PORT", "8973"))
LOG_FILE = OUTPUT_DIR / ".generation.log"
JOB_FILE = OUTPUT_DIR / ".generation.job.json"
QUEUE_FILE = OUTPUT_DIR / ".generation.queue.json"
MAX_CELLS = 400 # garde-fou : ~400 km² max par demande (zones dessinées)
MAX_CELLS_ALL = 10000 # garde-fou : passe globale « tout compléter » (dalles déjà présentes)
MAX_EXPORT_TILES = 64 # garde-fou : dalles assemblables par export (mémoire)
# Machine de traitement distante (webapp légère type Raspberry Pi) : si
# définie, la génération (téléchargement IGN + traitement) y est déléguée.
GENERATION_URL = (os.environ.get("LIDAR_GENERATION_URL") or "").rstrip("/")
# Jeton partagé : la webapp légère le présente à la machine de traitement
# (LIDAR_REMOTE_TOKEN ici, LIDAR_API_TOKEN là-bas). Optionnel.
REMOTE_TOKEN = os.environ.get("LIDAR_REMOTE_TOKEN") or None
API_TOKEN = os.environ.get("LIDAR_API_TOKEN") or None
# Réseau autorisé à lancer la (re)génération des tuiles (/api/generate) :
# liste de CIDR séparés par virgules, vide = restriction levée. Par défaut :
# boucle locale + plages privées RFC1918 — couvre le LAN, la machine hôte
# Docker (connexions via la passerelle 172.x) et localhost, tout en refusant
# les clients venus d'Internet (IP publiques non routables vers le privé).
# L'IP du client est l'IP de connexion (conservée par le DNAT Docker) ; si
# elle est dans le réseau autorisé, X-Forwarded-For est honoré (reverse proxy
# local qui expose la carte au-dehors — un client distant ne peut pas le
# forger puisque sa connexion directe est déjà hors réseau).
REGEN_CIDR = (os.environ.get(
"LIDAR_REGEN_CIDR",
"127.0.0.0/8,::1,10.0.0.0/8,172.16.0.0/12,192.168.0.0/16") or "").strip()
# Résolution générée par /api/generate (cf. _build_command) : 0,2 m seule,
# la 0,5 m n'est plus produite. La détection des tuiles à compléter exige
# cette même résolution : une tuile n'est complète que si elle possède les
# visualisations demandées en 0,2 m.
GENERATE_RESOLUTIONS = (0.2,)
app = FastAPI(title="Carte LiDAR — génération de zones")
# Visibilité de la configuration au démarrage (docker compose logs) : le
# backend de génération est TOUJOURS défini par variable d'environnement.
logger = logging.getLogger("lidar")
if not logger.handlers:
# Hors cli.py (uvicorn seul), le logger 'lidar' n'a aucun handler — en
# ajouter un pour que la config soit visible dans les logs du conteneur.
_h = logging.StreamHandler(sys.stdout)
_h.setFormatter(logging.Formatter("%(message)s"))
logger.addHandler(_h)
logger.setLevel(logging.INFO)
if GENERATION_URL:
logger.info(f"Backend de génération des tuiles : {GENERATION_URL} "
f"(LIDAR_GENERATION_URL)")
else:
logger.info("Backend de génération des tuiles : local — consultation "
"autonome du cache, sans LIDAR_GENERATION_URL")
@app.middleware("http")
async def _log_client_ip(request: Request, call_next):
xff = request.headers.get("x-forwarded-for", "").split(",")[0].strip()
peer = request.client.host if request.client else None
ip = xff or peer
path = request.url.path
# Ne logguer que les requêtes API (pas les fichiers statiques)
if path.startswith("/api/") or path == "/":
logger.info(f"{request.method} {path} — client : {ip or '?'}")
return await call_next(request)
# CSS/JS de l'interface : servis depuis l'IMAGE (bâchés au build, cf.
# Dockerfile / Dockerfile.webapp), plus depuis le volume — le code de
# l'interface suit l'image, pas le cache de tuiles. En repli (dev, source
# non bâchée), on les régénère depuis les constantes de l'index ; si le
# dossier du paquet n'est pas inscriptible (install système sans bake),
# repli dans un répertoire temporaire pour ne pas casser l'import.
_assets_dir = Path(__file__).resolve().parent / "webapp_assets"
if not (_assets_dir / "app.js").is_file():
from .index import _write_ui_assets
try:
_write_ui_assets(_assets_dir)
except OSError:
import tempfile
_assets_dir = Path(tempfile.gettempdir()) / "lidar_webapp_assets"
_write_ui_assets(_assets_dir)
# --- Cache à la demande : la visualisation peuple le cache local -------------
# Les images viennent de la machine de traitement par HTTP : une image absente
# du disque — ou plus ancienne que la version ?v= référencée (mtime ms de la
# source sur le worker) — est téléchargée depuis GENERATION_URL, écrite sur
# disque puis servie. Requêtes concurrentes dédupliquées (verrou par fichier),
# flux réseau plafonné (sémaphore).
from fastapi.staticfiles import StaticFiles as _StaticFiles
from starlette.concurrency import run_in_threadpool
_ONDEMAND_PREFIXES = ("visualisations", "index_thumbs", "index_subtiles")
_FETCH_SEMAPHORE = threading.Semaphore(4)
_FETCH_LOCKS = {}
_FETCH_LOCKS_GUARD = threading.Lock()
_MAX_CACHE_BYTES = 256 * 1024 * 1024 # garde-fou par fichier (dalle AVIF ~35 Mo)
# Disjoncteur : la machine de traitement n'est pas toujours allumée. Après
# quelques échecs consécutifs, les tentatives (index + images) sont suspendues
# 2 min — réponses locales instantanées, aucun martèlement ni erreur cliente.
_WORKER_HEALTH = {"fails": 0, "until": 0.0}
_WORKER_OFFLINE_AFTER = 3
_WORKER_OFFLINE_SECONDS = 120.0
def _worker_offline():
return time.time() < _WORKER_HEALTH["until"]
def _worker_mark(ok):
if ok:
_WORKER_HEALTH.update(fails=0, until=0.0)
else:
_WORKER_HEALTH["fails"] += 1
if (_WORKER_HEALTH["fails"] >= _WORKER_OFFLINE_AFTER
and _WORKER_HEALTH["until"] == 0.0):
_WORKER_HEALTH["until"] = time.time() + _WORKER_OFFLINE_SECONDS
logger.info("Machine de traitement injoignable — "
"rapatriement à la demande suspendu 2 min")
def _version_from_query(scope):
"""Version ?v= (mtime ms de la source) demandée par le client, ou None."""
try:
qs = urllib.parse.parse_qs(scope.get("query_string", b"").decode("ascii"))
v = qs.get("v", [None])[0]
return int(v) if v and v.isdigit() else None
except Exception:
return None
def _safe_rel_path(url_path):
"""Chemin relatif validé (préfixe autorisé, composants sains), ou None."""
rel = url_path.lstrip("/")
prefix = rel.split("/", 1)[0]
if prefix not in _ONDEMAND_PREFIXES:
return None
parts = [p for p in rel.split("/") if p not in ("", ".")]
if not parts or any(p == ".." or not re.fullmatch(r"[A-Za-z0-9_.-]+", p)
for p in parts):
return None
return "/".join(parts)
def _fetch_remote_to_cache(rel_path):
"""Télécharge GENERATION_URL/rel_path vers OUTPUT_DIR/rel_path (atomique)."""
if _worker_offline():
return False
dest = OUTPUT_DIR / rel_path
dest.parent.mkdir(parents=True, exist_ok=True)
tmp = dest.with_name(dest.name + ".part")
url = f"{GENERATION_URL}/{urllib.parse.quote(rel_path)}"
try:
with _FETCH_SEMAPHORE:
req = urllib.request.Request(url, headers={"User-Agent": "lidar-webapp-cache"})
with urllib.request.urlopen(req, timeout=60) as r, open(tmp, "wb") as f:
total = 0
while True:
chunk = r.read(1 << 20)
if not chunk:
break
total += len(chunk)
if total > _MAX_CACHE_BYTES:
raise IOError(f"fichier trop volumineux : {rel_path}")
f.write(chunk)
os.replace(tmp, dest)
_worker_mark(True)
logger.info(f"Cache à la demande : {rel_path} rapatrié "
f"({total / 1e6:.1f} Mo)")
return True
except Exception as e:
_worker_mark(False)
tmp.unlink(missing_ok=True)
logger.debug(f"Cache à la demande : échec {rel_path} : {e}")
return False
class _NoCacheStaticFiles(_StaticFiles):
"""Fichiers statiques servis sans cache navigateur (revalidation 304).
Les images (visualisations, vignettes, sous-tuiles, DTM) n'ont pas
toujours une URL invalidée après régénération : le ?v= suit la mtime de
la SOURCE, pas celle du fichier servi (vignette recalculée après coup,
cache webapp rapatrié…). Sans Cache-Control, le cache heuristique du
navigateur peut continuer d'afficher l'ancien rendu sur la même URL.
no-cache force la revalidation à chaque affichage — ETag/Last-Modified
la rend quasi gratuite quand le fichier n'a pas changé (réponse 304).
"""
async def get_response(self, path, scope):
resp = await super().get_response(path, scope)
resp.headers["Cache-Control"] = "no-cache, must-revalidate"
return resp
class _OnDemandStaticFiles(_NoCacheStaticFiles):
"""Montage statique qui peuple le cache depuis la machine de traitement.
Fichier local absent, ou périmé (?v= plus récent que sa mtime) → téléchar-
gement puis service ; sans GENERATION_URL ou en cas d'échec, comportement
statique normal (404, ou version locale périmée plutôt que rien).
url_prefix : préfixe du montage (ex. visualisations) — dans une app
montée, scope["path"] est le sous-chemin sans préfixe, on le reconstruit.
"""
def __init__(self, *, directory, url_prefix):
super().__init__(directory=directory)
self.url_prefix = url_prefix
async def get_response(self, path, scope):
want = None
if GENERATION_URL:
rel = _safe_rel_path(f"/{self.url_prefix}/{path}")
if rel is not None:
local = OUTPUT_DIR / rel
if not local.is_file():
want = rel
else:
v = _version_from_query(scope)
if v is not None and int(local.stat().st_mtime * 1000) + 500 < v:
want = rel # tuile régénérée sur le worker depuis le cache
if want is not None:
with _FETCH_LOCKS_GUARD:
lock = _FETCH_LOCKS.setdefault(want, threading.Lock())
# Téléchargement bloquant hors de la boucle d'événements ; le
# verrou par fichier déduplique les requêtes concurrentes.
await run_in_threadpool(self._fetch_under_lock, want, lock)
return await super().get_response(path, scope)
@staticmethod
def _fetch_under_lock(rel, lock):
with lock:
return _fetch_remote_to_cache(rel)
for name in ("index_thumbs", "index_subtiles", "visualisations", "DTM"):
_dir = OUTPUT_DIR / name
_dir.mkdir(parents=True, exist_ok=True)
if name in _ONDEMAND_PREFIXES and GENERATION_URL:
app.mount(f"/{name}",
_OnDemandStaticFiles(directory=str(_dir), url_prefix=name),
name=name)
else:
app.mount(f"/{name}", _NoCacheStaticFiles(directory=str(_dir)), name=name)
@app.get("/assets/{file_path:path}")
def assets(file_path: str):
"""Sert les fichiers de l'interface sans cache (régénérés à chaque rebuild)."""
base = _assets_dir.resolve()
p = (_assets_dir / file_path).resolve()
if base not in p.parents or not p.is_file():
raise HTTPException(404, f"asset introuvable : {file_path}")
return FileResponse(str(p), headers={"Cache-Control": "no-cache, must-revalidate"})
# Méthodes de classification du sol acceptées (mêmes valeurs que --ground-classification).
GROUND_CLASS_METHODS = ("auto", "ign", "smrf", "csf")
def _require_token(x_lidar_token: str = Header(None)):
"""Protège les routes mutantes si LIDAR_API_TOKEN est défini (sinon no-op).
Utilisé par la machine de traitement pour n'accepter que la webapp légère
autorisée (qui présente LIDAR_REMOTE_TOKEN) sur le réseau local.
"""
if API_TOKEN:
import hmac
presented = x_lidar_token or ""
if not hmac.compare_digest(presented, API_TOKEN):
raise HTTPException(401, "token d'API manquant ou invalide")
def _ip_in_regen_cidr(ip):
"""Vrai si l'IP appartient à un des CIDR autorisés (LIDAR_REGEN_CIDR).
LIDAR_REGEN_CIDR accepte plusieurs réseaux séparés par des virgules
(ex. « 192.168.1.0/24,10.0.0.0/8 ») ; les entrées invalides sont ignorées.
"""
if not REGEN_CIDR:
return True # restriction désactivée
if not ip:
return False
import ipaddress
try:
addr = ipaddress.ip_address(ip)
except ValueError:
return False
for part in REGEN_CIDR.split(","):
part = part.strip()
if not part:
continue
try:
if addr in ipaddress.ip_network(part, strict=False):
return True
except ValueError:
continue
return False
def _client_ip(request):
"""IP d'origine du client pour la restriction de génération.
Le pair direct doit être dans le réseau autorisé ; dans ce cas seulement,
X-Forwarded-For (reverse proxy local) désigne le client réel derrière lui.
Un client externe connecté en direct est refusé sur son IP de connexion,
quel que soit le X-Forwarded-For qu'il annonce.
"""
if request is None:
return None
peer = getattr(request, "client", None)
peer_ip = getattr(peer, "host", None) if peer else None
if not _ip_in_regen_cidr(peer_ip):
return peer_ip # hors réseau (ou inconnu) : jugé sur cette IP
xff = request.headers.get("x-forwarded-for") if hasattr(request, "headers") else None
if xff:
first = xff.split(",")[0].strip()
if first:
return first
return peer_ip
def _require_lan_for_generation(request: Request):
"""Réserve /api/generate aux clients du réseau local (LIDAR_REGEN_CIDR)."""
ip = _client_ip(request)
if not _ip_in_regen_cidr(ip):
raise HTTPException(
403, f"génération de tuiles réservée au réseau local"
f"{f' ({REGEN_CIDR})' if REGEN_CIDR else ''} — "
f"client : {ip or 'IP inconnue'}")
def _proxy_api(method, path, payload=None, timeout=30):
"""Transmet un appel d'API à la machine de traitement (mode webapp légère).
Retourne la réponse JSON distante. Les erreurs HTTP distantes (ex: 409 une
génération est déjà en cours) sont relayées telles quelles ; une machine
injoignable devient un 503 explicite pour l'interface.
"""
import urllib.error
import urllib.request
req = urllib.request.Request(
GENERATION_URL + path,
data=json.dumps(payload).encode("utf-8") if payload is not None else None,
method=method)
req.add_header("Content-Type", "application/json")
if REMOTE_TOKEN:
req.add_header("X-Lidar-Token", REMOTE_TOKEN)
try:
with urllib.request.urlopen(req, timeout=timeout) as resp:
return json.loads(resp.read().decode("utf-8"))
except urllib.error.HTTPError as e:
try:
detail = json.loads(e.read().decode("utf-8")).get("detail", str(e))
except (ValueError, OSError):
detail = str(e)
raise HTTPException(e.code, detail)
except (urllib.error.URLError, OSError) as e:
raise HTTPException(503, f"machine de traitement injoignable "
f"({GENERATION_URL}) : {e}")
def _viz_step_names():
"""Noms des visualisations acceptées (étapes VIZ_STEPS du pipeline).
Sur la webapp légère (sans dépendances de traitement), repli sur le
registre de l'index — le même ordre que VIZ_STEPS.
"""
try:
from .pipeline import VIZ_STEPS
return [name for name, _ in VIZ_STEPS]
except ImportError:
return _fallback_viz_step_names()
def _fallback_viz_step_names():
"""Noms d'étapes dérivés du registre léger de l'index (sans pipeline)."""
from .index import KEYWORD_TO_STEP, VIZ_LABELS
return [KEYWORD_TO_STEP.get(k, k) for k in VIZ_LABELS]
def _panel_viz_steps():
"""Couches réellement affichées par la webapp, en noms d'étapes --only.
Source unique : PANEL_VIZ (index.py) — le panneau de couches et le
sélecteur de génération sont pilotés par ce registre, les valeurs par
défaut de /api/preview et /api/generate doivent donc l'être aussi, sans
quoi une tuile générée « par défaut » serait incomplète pour le panneau.
"""
from .index import PANEL_VIZ, KEYWORD_TO_STEP
if not PANEL_VIZ:
return ["aspect"]
return [KEYWORD_TO_STEP.get(v, v) for v in PANEL_VIZ]
def _viz_step_labels():
"""Libellés français des étapes de visualisation (clé = nom d'étape).
VIZ_LABELS est indexée par mot-clé de nom de fichier ; trois étapes ont un
nom de sortie différent (cf. _expected_output_path dans le pipeline).
"""
from .index import VIZ_LABELS, step_to_keyword
return {name: VIZ_LABELS.get(step_to_keyword(name), name)
for name in _viz_step_names()}
class PreviewRequest(BaseModel):
bbox: Optional[list] = Field(None, description="[ouest, sud, est, nord] en WGS84 "
"(requis sauf all_missing=true)")
regenerate: bool = Field(False, description="Inclure les tuiles déjà générées")
viz: Optional[list] = Field(None,
description="Visualisations demandées, noms d'étapes du "
"pipeline (ex: aspect, wavelet) ; défaut : les "
"couches affichées dans le panneau. "
"Une tuile existante est considérée faite "
"seulement si elle les possède toutes")
all_missing: bool = Field(False,
description="Ignorer la bbox : toutes les dalles LHD "
"présentes dans input/ qui manquent au moins "
"une des visualisations demandées (sans "
"téléchargement)")
class GenerateRequest(BaseModel):
tiles: list = Field(..., description="liste [col, row] (entiers km L93)")
regenerate: bool = Field(False, description="Régénérer les tuiles déjà générées "
"(visualisations refaites, classification conservée)")
reclassify: bool = Field(False,
description="Relancer la classification du sol même si "
"la méthode demandée est déjà en cache "
"(--force-classification). Défaut : conserver "
"la classification existante ; choisir une "
"autre méthode reclassifie de toute façon "
"les tuiles concernées")
ground_class: str = Field("ign",
description="Méthode de classification du sol : "
"auto, ign, smrf, csf")
ign_classes: str = Field("sol",
description="Classes LAS pour le MNT IGN : liste noms ou "
"codes séparés par virgules — sol(2), "
"unclassified(1), eau(9), virtuel(66), "
"pont(17), sursol(64). Mode pur, "
"aucune retouche. (défaut: sol)")
bare_earth: bool = Field(False,
description="Sol nu : DTM au retour le plus bas de "
"chaque cellule (requalifie le point le plus "
"bas en terrain)")
viz: list = Field(None,
description="Visualisations à générer, noms d'étapes du "
"pipeline (ex: aspect, wavelet, slope) ; "
"défaut : les couches affichées dans le panneau")
all_missing: bool = Field(False,
description="Ignorer la liste de tuiles : traiter toutes "
"les dalles LHD présentes dans input/ qui "
"manquent au moins une des visualisations "
"demandées (sans téléchargement)")
class ExportRequest(BaseModel):
tiles: list = Field(..., description="liste [col, row] des dalles adjacentes "
"à assembler (entiers km L93)")
viz: list = Field(..., description="mots-clés de visualisation (ex: slope, "
"hillshade_multi) ; plusieurs possibles en PDF "
"(une page par visualisation), une seule pour une image")
resolution: float = Field(0.5, description="résolution des dalles en m/px "
"(0.5 ou 0.2)")
format: str = Field("jpeg", description="format de sortie : png, jpeg, "
"webp ou pdf")
max_side: int = Field(0, description="réduire le plus grand côté de la "
"mosaïque à N pixels (0 = pleine résolution) — "
"recommandé sur téléphone (4096)")
# --- État du job de génération -------------------------------------------
_job: dict = {"proc": None, "started": None, "returncode": None, "cmd": None, "finished": None, "qid": None, "run_id": None}
_job_lock = threading.Lock()
def _save_job_state():
"""Écrit l'état du job sur disque — survit au redémarrage du serveur.
Sans cela, un restart du conteneur perd tout : /api/status ne rapporte
plus le dernier run et la file de génération reste vide au rechargement.
Best-effort : n'acquiert pas le verrou (appelé aussi sous verrou).
"""
state = {k: _job.get(k) for k in ("started", "returncode", "cmd", "finished", "qid", "run_id")}
try:
JOB_FILE.write_text(json.dumps(state), encoding="utf-8")
except OSError:
pass
def _load_job_state():
"""Recharge l'état du dernier job au démarrage du serveur."""
try:
state = json.loads(JOB_FILE.read_text(encoding="utf-8"))
except (OSError, ValueError):
return
with _job_lock:
_job["started"] = state.get("started")
_job["returncode"] = state.get("returncode")
_job["cmd"] = state.get("cmd")
_job["finished"] = state.get("finished")
_job["qid"] = state.get("qid")
_job["run_id"] = state.get("run_id")
_load_job_state()
# --- File d'attente des demandes de génération -----------------------------
# Une demande arrivée pendant un run n'est plus rejetée (ancien 409) : elle
# part en file et démarre automatiquement à la fin du run en cours — le
# travail déjà lancé n'est jamais coupé. La file survit aux redémarrages du
# serveur (fichier .generation.queue.json) et est vidable via /api/queue/clear.
_queue = []
_queue_seq = 0
def _save_queue():
"""Écrit la file d'attente sur disque (atomique, best-effort)."""
try:
tmp = QUEUE_FILE.with_suffix(".tmp")
tmp.write_text(json.dumps(_queue), encoding="utf-8")
os.replace(tmp, QUEUE_FILE)
except OSError:
pass
def _load_queue():
"""Recharge la file d'attente au démarrage (et son compteur d'ids)."""
global _queue_seq
try:
data = json.loads(QUEUE_FILE.read_text(encoding="utf-8"))
if isinstance(data, list):
_queue[:] = [it for it in data
if isinstance(it, dict) and isinstance(it.get("req"), dict)]
except (OSError, ValueError):
pass
_queue_seq = max((it.get("id", 0) for it in _queue), default=0)
def _queue_summary():
"""Aperçu de la file pour /api/status (sans la requête brute)."""
return [{"id": it.get("id"),
"tuiles": len(it["req"].get("tiles") or []),
"toutes": bool(it["req"].get("all_missing")),
"en_attente_depuis": it.get("queued_at")}
for it in _queue]
_load_queue()
# --- Presets de couches (partagés entre navigateurs) --------------------------
# Stockés dans OUTPUT_DIR/.presets.json : un preset est un jeu nommé de
# couches visibles + opacités + ordre de pile + fond OSM, créé depuis
# l'interface (bouton 💾) et partagé par tous les navigateurs connectés à
# cette webapp (y compris le Pi en mode deux machines). Le localStorage du
# navigateur reste le repli quand l'interface est servie sans webapp.
PRESETS_FILE = OUTPUT_DIR / ".presets.json"
_presets_lock = threading.Lock()
def _load_presets():
"""Lis la liste des presets depuis disque (vide si absent/corrompu)."""
try:
data = json.loads(PRESETS_FILE.read_text(encoding="utf-8"))
if isinstance(data, list):
return data
except (OSError, ValueError):
pass
return []
def _save_presets(presets):
"""Écrit la liste des presets atomiquement (fichier temporaire + replace)."""
tmp = PRESETS_FILE.with_suffix(".tmp")
tmp.write_text(json.dumps(presets, ensure_ascii=False), encoding="utf-8")
os.replace(tmp, PRESETS_FILE)
class PresetRequest(BaseModel):
id: str = Field(..., description="identifiant du preset (attribué par l'interface)")
label: str = Field(..., description="nom affiché dans l'interface")
on: list = Field(default_factory=list, description="couches visibles (clés, ordre de pile)")
opacity: dict = Field(default_factory=dict, description="opacité de chaque couche, 0–1")
osm: dict = Field(default_factory=dict, description="fond de carte OSM : on, opacity, dark")
def _sanitize_preset(req: PresetRequest) -> dict:
"""Preset à stocker : label assaini, opacités bornées entre 0 et 1."""
on = [str(k) for k in (req.on or []) if str(k).strip()]
opacity = {}
for k, v in (req.opacity or {}).items():
try:
opacity[str(k)] = min(1.0, max(0.0, float(v)))
except (TypeError, ValueError):
continue
osm_in = req.osm if isinstance(req.osm, dict) else {}
osm = {"on": bool(osm_in.get("on", False))}
try:
osm["opacity"] = min(1.0, max(0.0, float(osm_in["opacity"])))
except (KeyError, TypeError, ValueError):
pass
if "dark" in osm_in:
osm["dark"] = bool(osm_in["dark"])
pid = re.sub(r"[^a-z0-9_-]", "", str(req.id).lower())
return {
"id": pid or "p" + str(int(time.time() * 1000)),
"label": (str(req.label).strip() or "Preset")[:80],
"on": on,
"opacity": opacity,
"osm": osm,
}
@app.get("/api/presets")
def list_presets():
"""Presets de couches personnalisés (partagés entre tous les navigateurs)."""
with _presets_lock:
return {"presets": _load_presets()}
@app.post("/api/presets", dependencies=[Depends(_require_token)])
def upsert_preset(req: PresetRequest):
"""Crée ou remplace un preset de couches (même identifiant = remplacement)."""
p = _sanitize_preset(req)
with _presets_lock:
presets = [x for x in _load_presets() if x.get("id") != p["id"]]
presets.append(p)
_save_presets(presets)
return p
@app.delete("/api/presets/{preset_id}", dependencies=[Depends(_require_token)])
def delete_preset(preset_id: str):
"""Supprime un preset de couches (404 s'il est inconnu)."""
with _presets_lock:
presets = _load_presets()
kept = [x for x in presets if x.get("id") != preset_id]
if len(kept) == len(presets):
raise HTTPException(404, "preset inconnu")
_save_presets(kept)
return {"ok": True}
def bbox_to_cells(w, s, e, n):
"""Cellules L93 de 1 km (col, row) intersectant une bbox WGS84.
Une cellule (col, row) couvre X ∈ [col, col+1] km, Y ∈ [row-1, row] km.
Conversion via rasterio/PROJ si disponible, sinon pyproj, sinon
l'approximation affine de l'index (webapp légère sans GDAL).
"""
from .index import _approx_wgs84_to_l93
try:
from rasterio.warp import transform as warp_transform
xs, ys = warp_transform("EPSG:4326", "EPSG:2154", [w, e, w, e], [s, s, n, n])
except ImportError:
try:
from pyproj import Transformer
transformer = Transformer.from_crs("EPSG:4326", "EPSG:2154", always_xy=True)
xs, ys = transformer.transform([w, e, w, e], [s, s, n, n])
except ImportError:
pts = [_approx_wgs84_to_l93(lon, lat)
for lon, lat in ((w, s), (e, s), (w, n), (e, n))]
xs = [p[0] for p in pts]
ys = [p[1] for p in pts]
# transform renvoie (xs, ys) dans la CRS cible
min_x, max_x = min(xs) + 0.5, max(xs) - 0.5 # rétrécit d'1 m : bords exclus
min_y, max_y = min(ys) + 0.5, max(ys) - 0.5
if max_x <= min_x or max_y <= min_y:
return []
cols = range(int(math.floor(min_x / 1000)), int(math.floor(max_x / 1000)) + 1)
rows = range(int(math.floor(min_y / 1000)) + 1, int(math.floor(max_y / 1000)) + 2)
return [(c, r) for r in rows for c in cols]
def point_to_cell(lon, lat):
"""Cellule L93 de 1 km (col, row) contenant un point WGS84.
Cohérent avec bbox_to_cells : une cellule (col, row) couvre
X ∈ [col, col+1] km, Y ∈ [row-1, row] km — sert au clic carte pour
sélectionner une dalle précise, même sans données existantes.
"""
from .index import _approx_wgs84_to_l93
try:
from rasterio.warp import transform as warp_transform
x, y = warp_transform("EPSG:4326", "EPSG:2154", [lon], [lat])
x, y = x[0], y[0]
except ImportError:
try:
from pyproj import Transformer
transformer = Transformer.from_crs("EPSG:4326", "EPSG:2154", always_xy=True)
x, y = transformer.transform(lon, lat)
except ImportError:
x, y = _approx_wgs84_to_l93(lon, lat)
return int(math.floor(x / 1000)), int(math.floor(y / 1000)) + 1
def processed_cells(output_dir):
"""Ensemble des cellules (col, row) ayant déjà des visualisations."""
from .index import scan_tiles
tiles = scan_tiles(Path(output_dir) / "visualisations")
return {(t["col"], t["row"]) for t in tiles}
def laz_cells(input_dir):
"""Coordonnées (col, row) des dalles LHD présentes dans input/.
Repère les fichiers dont le nom est LHD_FXX_{col}_{row}_... (extensions
.copc.laz/.copc.las/.laz/.las) ; les autres fichiers sont ignorés.
"""
from .index import parse_basename_coords
cells = set()
for f in list(Path(input_dir).glob("*.laz")) + list(Path(input_dir).glob("*.las")):
name = f.name
for ext in ('.copc.laz', '.copc.las', '.laz', '.las'):
if name.lower().endswith(ext):
name = name[:-len(ext)]
break
coords = parse_basename_coords(name)
if coords:
cells.add(coords)
return sorted(cells)
def complete_cells(output_dir, viz_steps):
"""Cellules ayant TOUTES les visualisations demandées à la résolution de
génération (0,2 m).
viz_steps : noms d'étapes du pipeline (ex: 'aspect', 'pos_open'), convertis
en mots-clés de fichiers de sortie avant comparaison avec le disque.
"""
from .index import cells_with_all_viz, step_to_keyword
keys = [step_to_keyword(v) for v in viz_steps]
return cells_with_all_viz(Path(output_dir) / "visualisations", keys,
GENERATE_RESOLUTIONS)
def missing_cells_with_corners(cells, output_dir, include_done=False, viz=None):
"""Filtre les cellules à traiter et calcule leurs coins WGS84.
Sans viz : une cellule ayant des visualisations est considérée faite.
Avec viz (noms d'étapes) : seules les cellules disposant de TOUTES les
visualisations demandées (aux résolutions de la génération) sont faites —
les tuiles existantes mais incomplètes restent incluses, afin de générer
les visualisations manquantes sans --force (le pipeline ignore l'existant).
Retourne [{col, row, corners: [[lat, lon] × 4 SW,SE,NE,NW}].
"""
from .index import attach_gps_bounds
if viz:
done = complete_cells(output_dir, viz)
else:
done = processed_cells(output_dir)
todo = [{"col": c, "row": r} for (c, r) in cells
if include_done or (c, r) not in done]
if todo:
attach_gps_bounds(todo)
return todo
@app.get("/")
def root():
index = OUTPUT_DIR / "index.html"
if not index.exists():
return JSONResponse({"erreur": "index.html introuvable — lancez d'abord le pipeline"},
status_code=404)
# index.html est régénéré à chaque passe du pipeline : on interdit le cache
# navigateur pour ne pas servir une version périmée (ex. menu de génération).
return FileResponse(
str(index), media_type="text/html",
headers={
"Cache-Control": "no-cache, no-store, must-revalidate",
"Pragma": "no-cache",
"Expires": "0",
})
@app.get("/api/status")
def status(request: Request = None):
# Droit de lancer une génération depuis cette IP (guide l'interface)
regen_allowed = _ip_in_regen_cidr(_client_ip(request))
if GENERATION_URL:
# Webapp légère : l'état vient de la machine de traitement ; le droit
# de génération reste jugé sur l'IP du navigateur (localement). Hors
# ligne : running=null (l'interface ignore, aucun faux « terminé »).
try:
data = _proxy_api("GET", "/api/status", timeout=10)
except HTTPException as e:
if e.status_code != 503:
raise
return {"running": None, "distant": True, "worker_offline": True,
"regen_allowed": regen_allowed}
data["distant"] = True
data["regen_allowed"] = regen_allowed
return data
from .progress import progress_snapshot
with _job_lock:
pending = bool(_queue)
if pending:
# File d'attente non vide et serveur libre (ex. redémarrage du
# conteneur en cours de file) : la demande suivante repart.
_start_next_queued()
with _job_lock:
proc = _job["proc"]
running = proc is not None and proc.poll() is None
return {
"running": running,
"started": _job["started"],
"returncode": _job["returncode"],
"cmd": _job["cmd"],
"qid": _job.get("qid"),
"queue": _queue_summary(),
"log": _tail_log(40),
# L'interface masque les boutons de génération hors réseau local
"regen_allowed": regen_allowed,
# Tuiles une par une (nom court, type, état) — agrégées depuis les
# événements écrits par le pipeline (.generation.events.jsonl)
"tiles": progress_snapshot(OUTPUT_DIR, viz_labels=_viz_step_labels(),
run_id=_job.get("run_id")),
}
@app.get("/api/layers")
def available_layers():
"""Couches proposées dans le panneau (clé → label), sans rebuild de l'index.
Toutes les couches connues du registre (VIZ_LABELS) détectées sur disque.
Permet à la carte de détecter celles apparues depuis le dernier
build_index et de proposer une actualisation.
"""
from .index import VIZ_LABELS
vis = OUTPUT_DIR / "visualisations"
found = {}
if vis.is_dir():
for key in VIZ_LABELS:
for ext in ("avif", "webp"):
if next(vis.glob(f"*/*_{key}.{ext}"), None) is not None:
found[key] = VIZ_LABELS.get(key, key)
break
return found
_REMOTE_INDEX = {"data": None, "fetched": 0.0}
_REMOTE_INDEX_TTL = 60.0
def _tile_entry_key(t):
"""Clé d'unicité d'une tuile affichée (position, résolution, quadrant)."""
return (t.get("col"), t.get("row"), t.get("resolution"),
t.get("sub_i"), t.get("sub_j"))
def _remote_tiles_data():
"""Index de la machine de traitement (cache 60 s), ou None si injoignable.
Une panne ne remonte pas : dernière réponse connue servie telle quelle
(stale-while-error), et le disjoncteur suspend les tentatives.
"""
if not GENERATION_URL or _worker_offline():
return _REMOTE_INDEX["data"]
now = time.time()
if _REMOTE_INDEX["fetched"] > now - _REMOTE_INDEX_TTL:
return _REMOTE_INDEX["data"]
try:
with urllib.request.urlopen(f"{GENERATION_URL}/api/tiles", timeout=4) as r:
data = json.loads(r.read().decode("utf-8"))
_worker_mark(True)
_REMOTE_INDEX.update({"data": data, "fetched": now})
except Exception as e:
_worker_mark(False)
_REMOTE_INDEX["fetched"] = now # ne pas marteler à chaque requête
logger.debug(f"Index distant injoignable : {e}")
return _REMOTE_INDEX["data"]
@app.get("/api/tiles")
def tiles_data(stamp: Optional[float] = None):
"""Données de carte — cache local + index de la machine de traitement.
Le pipeline lancé par /api/generate tourne avec --incremental-index : il
réécrit index_tiles.json après chaque tuile terminée. La carte sonde cet
endpoint et fusionne les nouveautés. Sur la webapp légère, l'index du
worker (cache 60 s) fait autorité : les tuiles absentes du cache local y
figurent quand même — leurs images sont rapatriées à la visualisation par
le montage statique à la demande. Worker hors ligne : cache local seul,
sans erreur. Avec ?stamp=X : réponse allégée ({tiles: null}) si rien n'a
changé (le stamp combine les deux index).
"""
path = OUTPUT_DIR / "index_tiles.json"
try:
local_stamp = path.stat().st_mtime
data = json.loads(path.read_text(encoding="utf-8"))
except (OSError, ValueError):
local_stamp = None
data = None
remote = _remote_tiles_data()
remote_stamp = (remote or {}).get("stamp")
stamps = [s for s in (local_stamp, remote_stamp) if s is not None]
current = max(stamps) if stamps else None
if stamp is not None and current is not None and abs(current - stamp) < 1e-4:
return {"stamp": current, "tiles": None}
if data is None and remote is None:
return {"stamp": None, "tiles": None}
# L'index distant prime (tuiles régénérées : ?v= neuf → re-téléchargement) ;
# le local complète les positions que le worker n'a plus (historique).
merged = list((remote or {}).get("tiles") or [])
keys = {_tile_entry_key(t) for t in merged}
for t in (data or {}).get("tiles") or []:
k = _tile_entry_key(t)
if k not in keys:
merged.append(t)
keys.add(k)
viz_meta = dict((data or {}).get("viz_meta") or {})
viz_meta.update((remote or {}).get("viz_meta") or {})
# Le panneau de couches est figé (PANEL_VIZ) : un index obsolète (ex.
# worker antérieur à la restriction) ne doit pas réintroduire de couches
# retirées dans le menu au rechargement de la page.
from .index import PANEL_VIZ
if PANEL_VIZ is not None:
allowed = set(PANEL_VIZ)
viz_meta = {k: v for k, v in viz_meta.items() if k in allowed}
return {
"tiles": merged,
"viz_meta": viz_meta,
"stats": (remote or data or {}).get("stats") or {},
"stamp": current,
}
# --- Rebuild de l'index en arrière-plan --------------------------------------
# Les tuiles sont servies à la demande (cache peuplé depuis la machine de
# traitement) : /api/sync ne fait plus qu'en régénérer l'index local
# (vignettes + carte).
_rebuild = {"running": False, "error": None, "done": None, "phase": None}
def _merge_remote_index():
"""Fusionne l'index du worker dans index_tiles.json et la coquille HTML.
Le rebuild local ne voit que le cache : sans fusion, la coquille embarque
un sous-ensemble des tuiles (celles déjà rapatriées) et le premier
chargement de la carte n'affiche pas les tuiles récentes du worker —
elles n'arrivent qu'au sondage /api/tiles qui suit. Mêmes règles que
tiles_data : l'index distant prime (tuiles régénérées), le local complète
les positions que le worker n'a plus (historique).
"""
remote = _remote_tiles_data() # frais : _start_rebuild a vidé le cache
if not remote or not remote.get("tiles"):
return
path = OUTPUT_DIR / "index_tiles.json"
try:
data = json.loads(path.read_text(encoding="utf-8"))
except (OSError, ValueError):
data = {}
merged = list(remote["tiles"])
keys = {_tile_entry_key(t) for t in merged}
for t in data.get("tiles") or []:
k = _tile_entry_key(t)
if k not in keys:
merged.append(t)
keys.add(k)
viz_meta = dict(data.get("viz_meta") or {})
viz_meta.update(remote.get("viz_meta") or {})
from .index import PANEL_VIZ
if PANEL_VIZ is not None:
allowed = set(PANEL_VIZ)
viz_meta = {k: v for k, v in viz_meta.items() if k in allowed}
path.write_text(json.dumps({
"tiles": merged,
"viz_meta": viz_meta,
# Les stats décrivent les réglages de l'INTERFACE (couches par
# défaut) : elles viennent du rebuild local, fait avec le code de
# cette webapp — celles du worker ne servent que de repli tant
# qu'aucun index local n'existe (sinon une machine mise à jour avant
# l'autre afficherait les vieux réglages du worker).
"stats": data.get("stats") or remote.get("stats") or {},
}, ensure_ascii=False), encoding="utf-8")
# Coquille HTML : la liste embarquée doit être la même (fusionnée). Le
# JSON ne contient jamais de « ; » : la première séquence « ]; » après
# « const TILES = [ » termine forcément l'instruction.
html_path = OUTPUT_DIR / "index.html"
try:
html = html_path.read_text(encoding="utf-8")
new_html = re.sub(r"const TILES = \[.*?\];",
"const TILES = " + json.dumps(merged, ensure_ascii=False)
.replace("</", "<\\/") + ";",
html, count=1, flags=re.S)
if new_html != html:
html_path.write_text(new_html, encoding="utf-8")
except OSError:
pass
logger.info(f"Index fusionné avec la machine de traitement : "
f"{len(merged)} tuile(s) au total")
def _start_rebuild():
"""Lance en arrière-plan le rebuild de l'index (vignettes + carte).
L'état running est posé de façon SYNCHRONE, avant le fil : un GET
/api/sync juste après le POST doit voir le rebuild en cours — sinon il
lirait le done du rebuild précédent et l'interface rechargerait une
carte périmée. Le timestamp started retourné permet au client
d'attendre la fin de CE rebuild précis (done >= started).
"""
if _rebuild["running"]:
raise HTTPException(409, "un rebuild est déjà en cours")
started = time.time()
_rebuild.update({"running": True, "error": None, "done": None,
"phase": "index"})
# Index distant (mode deux machines) : cache 60 s vidé pour que le
# premier /api/tiles suivant reflète l'état du worker sans attendre.
_REMOTE_INDEX["fetched"] = 0.0
def _run():
try:
from .index import build_index
build_index(OUTPUT_DIR)
# Mode deux machines : le rebuild local ne voit que le cache —
# on y fusionne l'index du worker pour que coquille embarquée et
# /api/tiles servent la même vue complète dès le 1er chargement.
if GENERATION_URL:
try:
_merge_remote_index()
except Exception as e: # noqa: BLE001 — index local valable
logger.warning(f"Fusion de l'index distant ignorée : {e}")
except Exception as e: # noqa: BLE001 — remonté à l'UI via l'état
_rebuild["error"] = str(e)
finally:
_rebuild["running"] = False
_rebuild["phase"] = None
_rebuild["done"] = time.time()
threading.Thread(target=_run, daemon=True).start()
return {"demarré": True, "started": started}
@app.post("/api/sync", dependencies=[Depends(_require_token)])
def sync_and_rebuild():
"""Régénère l'index local (vignettes + carte) en arrière-plan.
Les tuiles étant servies à la demande, il n'y a plus de synchronisation à
exécuter : /api/sync se réduit à un rebuild de l'index (même
comportement que /api/rebuild).
"""
return _start_rebuild()
@app.get("/api/sync")
def sync_status():
return rebuild_status()
@app.post("/api/rebuild", dependencies=[Depends(_require_token)])
def rebuild_index():
"""Régénère la carte (index.html, vignettes, sous-tuiles) en arrière-plan."""
return _start_rebuild()
@app.get("/api/rebuild")
def rebuild_status():
return {"running": _rebuild["running"], "error": _rebuild["error"],
"done": _rebuild["done"], "phase": _rebuild["phase"]}
def _tail_log(n_lines):
"""Dernières lignes du journal — lecture depuis la fin uniquement.
Le journal peut atteindre plusieurs Mo et /api/status est sondé toutes
les ~2 s par l'interface : relire tout le fichier à chaque sondage
coûte cher, surtout sur Raspberry Pi.
"""
try:
with open(LOG_FILE, "rb") as fh:
fh.seek(0, os.SEEK_END)
size = fh.tell()
fh.seek(max(0, size - 64 * 1024))
chunk = fh.read().decode("utf-8", errors="replace")
return chunk.splitlines()[-n_lines:]
except OSError:
return []
@app.post("/api/preview", dependencies=[Depends(_require_token)])
def preview(req: PreviewRequest):
if GENERATION_URL:
# Webapp légère : le distant connaît input/ et output/ complets
return _proxy_api("POST", "/api/preview", json.loads(req.json()))
names = _viz_step_names()
viz = [v for v in (req.viz or []) if v] or _panel_viz_steps()
invalid = [v for v in viz if v not in names]
if invalid:
raise HTTPException(
400, f"visualisation invalide : {', '.join(invalid)} "
f"(attendues : {', '.join(names)})")
if req.all_missing:
# Passe globale : dalles LHD déjà présentes dans input/ (aucun téléchargement)
cells = laz_cells(INPUT_DIR)
capped = len(cells) > MAX_CELLS_ALL
todo = missing_cells_with_corners(cells[:MAX_CELLS_ALL], OUTPUT_DIR,
include_done=req.regenerate, viz=viz)
else:
if not req.bbox or len(req.bbox) != 4:
raise HTTPException(400, "bbox attendue : [ouest, sud, est, nord]")
w, s, e, n = (float(v) for v in req.bbox)
cells = bbox_to_cells(w, s, e, n)
capped = len(cells) > MAX_CELLS
todo = missing_cells_with_corners(cells[:MAX_CELLS], OUTPUT_DIR,
include_done=req.regenerate, viz=viz)
return {"count": len(todo), "capped": capped, "cells": todo}
@app.get("/api/cell")
def cell_at_point(lat: float, lng: float):
"""Dalle L93 de 1 km contenant un point WGS84 (clic sur la carte).
Retourne col, row et les coins GPS de la dalle — permet de sélectionner
une tuile individuelle même sans données existantes (téléchargement +
génération à la validation).
"""
if GENERATION_URL:
# Webapp légère : la conversion précise vit sur la machine de traitement
return _proxy_api("GET", f"/api/cell?lat={lat}&lng={lng}")
from .index import attach_gps_bounds
col, row = point_to_cell(lng, lat)
tile = {"col": col, "row": row}
attach_gps_bounds([tile])
return {"col": col, "row": row, "corners": tile.get("corners") or []}
def _build_command(tiles, regenerate=False, ground_class="ign", bare_earth=False, ign_classes="sol", viz=None, reclassify=False):
"""Commande de génération : téléchargement IGN + traitement des fichiers.
La classification du sol est choisie via `ground_class` (défaut : "ign",
pré-classification IGN ; le pipeline bascule sur SMRF si un fichier ne la
contient pas). Avec ign_classes, on choisit les classes LAS extraites
pour le MNT (mode pur, ex. "sol,unclassified" pour combler les trous
sans retouche). Avec regenerate=True, les visualisations des tuiles déjà
présentes sont refaites (--force) mais leur classification est conservée :
le pipeline réutilise le DTM en cache quand la méthode ne change pas.
Avec reclassify=True, --force-classification relance la classification
même à méthode inchangée (sinon, choisir une méthode différente via
ground_class/ign_classes reclassifie déjà les tuiles concernées).
Avec bare_earth=True, le DTM est ramené au retour le plus bas de chaque
cellule (sol nu). Avec viz, on choisit les visualisations générées
(noms d'étapes du pipeline, ex. ["aspect", "wavelet", "slope"] ;
défaut : aspect).
"""
from .fetch_ign import tile_filename
# -u : sortie non bufferisée — le journal .generation.log doit être
# lu en temps réel par /api/status (progression affichée dans l'UI).
cmd = [sys.executable, "-u", "-m", "lidar_pipeline", str(INPUT_DIR),
"-o", str(OUTPUT_DIR),
"-r", ",".join(str(r) for r in GENERATE_RESOLUTIONS),
"--only", *(viz or _panel_viz_steps()),
"--ground-classification", ground_class,
"--ign-classes", ign_classes]
if bare_earth:
cmd += ["--bare-earth"]
if regenerate:
cmd += ["--force"]
if reclassify:
cmd += ["--force-classification"]
if os.environ.get("LIDAR_GPU", "") == "1":
cmd += ["-g", "all", "-w", os.environ.get("LIDAR_WORKERS", "2")]
# Carte régénérée après chaque tuile terminée : la webapp l'affiche en
# direct via /api/tiles pendant le run.
cmd += ["--incremental-index"]
cmd += ["--fetch-tiles"]
cmd += [f"{c},{r}" for (c, r) in tiles]
cmd += ["--file"]
cmd += [tile_filename(c, r) for (c, r) in tiles]
return cmd
def _resolve_request(req):
"""Valide une demande de génération et résout sa liste de tuiles.
Retourne (tuiles [(col, row)], viz [noms d'étapes]). Réexécutée au
démarrage de chaque demande sortie de file : les dalles présentes dans
input/ peuvent avoir changé entre la mise en file et le départ du run.
"""
names = _viz_step_names()
viz = [v for v in (req.viz or []) if v] or _panel_viz_steps()
invalid = [v for v in viz if v not in names]
if invalid:
raise HTTPException(
400, f"visualisation invalide : {', '.join(invalid)} "
f"(attendues : {', '.join(names)})")
if req.all_missing:
# Passe globale : les dalles présentes dans input/ qui manquent les
# visualisations demandées (les complètes incluses seulement si regenerate)
done = set() if req.regenerate else complete_cells(OUTPUT_DIR, viz)
tiles = [c for c in laz_cells(INPUT_DIR)
if c not in done or req.regenerate]
if not tiles:
raise HTTPException(400, "aucune tuile à compléter dans input/")
# Même garde-fou que /api/preview : un run parti sur tout input/ sans
# plafond peut durer des jours sans moyen de l'interrompre proprement
if len(tiles) > MAX_CELLS_ALL:
raise HTTPException(
400, f"trop de tuiles à compléter ({len(tiles)}) — max "
f"{MAX_CELLS_ALL} par run ; relancez la passe globale à "
f"la fin de celui-ci")
else:
tiles = []
for pair in req.tiles:
if not (isinstance(pair, list) and len(pair) == 2):
raise HTTPException(400, f"tuile invalide : {pair!r} (attendu [col, row])")
tiles.append((int(pair[0]), int(pair[1])))
if not tiles:
raise HTTPException(400, "aucune tuile fournie")
if len(tiles) > MAX_CELLS:
raise HTTPException(400, f"trop de tuiles ({len(tiles)}) — max {MAX_CELLS}")
if req.ground_class not in GROUND_CLASS_METHODS:
raise HTTPException(
400, f"méthode de classification invalide : {req.ground_class!r} "
f"(attendu : {', '.join(GROUND_CLASS_METHODS)})")
return tiles, viz
def _launch_job(tiles, viz, req, qid=None):
"""Démarre un run du pipeline. À appeler sous _job_lock, serveur libre.
Retourne la commande lancée. qid : identifiant de file d'attente quand
la demande vient de la file (suivi côté interface).
"""
cmd = _build_command(tiles, regenerate=req.regenerate, ground_class=req.ground_class,
bare_earth=req.bare_earth, ign_classes=req.ign_classes,
viz=viz, reclassify=req.reclassify)
LOG_FILE.parent.mkdir(parents=True, exist_ok=True)
log_fh = open(LOG_FILE, "w", encoding="utf-8")
# Nouveau run : journal d'événements remis à zéro (les tuiles affichées
# dans la file correspondent au run qui démarre, pas au précédent) et
# identifiant de run transmis au pipeline (LIDAR_RUN_ID) : les événements
# tardifs d'un run annulé (workers survivants) sont filtrés à la lecture.
from .progress import reset_events
reset_events(OUTPUT_DIR)
run_id = uuid.uuid4().hex[:12]
_job.update({"proc": None, "started": time.time(), "returncode": None,
"cmd": cmd, "finished": None, "qid": qid, "run_id": run_id})
_save_job_state()
# start_new_session : le pipeline et ses workers/PDAL forment leur
# propre groupe de processus — /api/stop peut le tuer en bloc sans
# toucher la webapp, et le gestionnaire SIGTERM du pipeline (killpg)
# reste confiné à son groupe.
p = subprocess.Popen(cmd, stdout=log_fh, stderr=subprocess.STDOUT,
cwd="/app" if Path("/app").exists() else None,
env=dict(os.environ, LIDAR_RUN_ID=run_id),
start_new_session=True)
# proc posé AVANT le fil de veille : un run qui meurt instantanément
# déclencherait _watch → _start_next_queued avec _job["proc"] encore à
# None (serveur « libre ») et un second run partirait en parallèle.
_job["proc"] = p
def _watch():
rc = p.wait()
log_fh.close()
with _job_lock:
_job["returncode"] = rc
_job["finished"] = time.time()
_save_job_state()
# Le run est fini : la demande suivante de la file peut partir.
_start_next_queued()
threading.Thread(target=_watch, daemon=True).start()
return cmd
def _start_next_queued():
"""Démarre la première demande de la file dès que le serveur est libre.
Les demandes devenues invalides (ex. plus aucune tuile à compléter)
sont ignorées ; une demande repoussée en tête de file si un run part
entre-temps. Ne coupe jamais le run en cours.
"""
while True:
with _job_lock:
proc = _job["proc"]
if proc is not None and proc.poll() is None:
return False
if not _queue:
return False
item = _queue.pop(0)
_save_queue()
try:
req = GenerateRequest(**item["req"])
tiles, viz = _resolve_request(req)
except HTTPException as e:
logger.warning("Demande en file ignorée : %s", e.detail)
continue
with _job_lock:
proc = _job["proc"]
if proc is not None and proc.poll() is None:
_queue.insert(0, item) # un autre run est parti : remise en tête
_save_queue()
return False
_launch_job(tiles, viz, req, qid=item.get("id"))
return True
@app.post("/api/generate", dependencies=[Depends(_require_token),
Depends(_require_lan_for_generation)])
def generate(req: GenerateRequest):
if GENERATION_URL:
# Webapp légère : téléchargement + traitement sur la machine distante,
# qui valide les paramètres, met en file et renvoie son propre état.
data = _proxy_api("POST", "/api/generate", json.loads(req.json()))
data["distant"] = True
return data
tiles, viz = _resolve_request(req)
with _job_lock:
proc = _job["proc"]
running = proc is not None and proc.poll() is None
if running:
# Un run est en cours : la demande part en file d'attente et
# démarrera à sa fin — jamais de coupure du travail en place.
global _queue_seq
_queue_seq += 1
item = {"id": _queue_seq, "req": json.loads(req.json()),
"queued_at": time.time()}
_queue.append(item)
_save_queue()
return {"demarré": False, "en_file": len(_queue), "qid": item["id"],
"tuiles": len(tiles)}
cmd = _launch_job(tiles, viz, req)
return {"demarré": True, "tuiles": len(tiles), "commande": " ".join(cmd)}
@app.post("/api/queue/clear", dependencies=[Depends(_require_token),
Depends(_require_lan_for_generation)])
def queue_clear():
"""Retire les demandes en attente (le run en cours n'est pas touché)."""
if GENERATION_URL:
# Webapp légère : la file vit sur la machine de traitement
return _proxy_api("POST", "/api/queue/clear")
with _job_lock:
n = len(_queue)
_queue.clear()
_save_queue()
return {"retirées": n}
@app.post("/api/stop", dependencies=[Depends(_require_token),
Depends(_require_lan_for_generation)])
def stop_generation():
"""Arrête la génération en cours.
SIGTERM au pipeline : son gestionnaire de signal nettoie ses workers et
ses processus PDAL (killpg sur son propre groupe, cf. cli.py). Si le
processus ne meurt pas dans les 15 s (worker bloqué), SIGKILL du groupe
entier. L'état final (returncode négatif) est enregistré par le fil de
surveillance du job.
"""
if GENERATION_URL:
# Webapp légère : l'arrêt concerne le pipeline de la machine distante.
return _proxy_api("POST", "/api/stop")
with _job_lock:
proc = _job["proc"]
if proc is None or proc.poll() is not None:
raise HTTPException(409, "aucune génération en cours")
try:
proc.terminate() # SIGTERM → nettoyage propre du pipeline
except OSError:
pass
def _escalate():
try:
proc.wait(timeout=15)
except subprocess.TimeoutExpired:
try:
os.killpg(proc.pid, signal.SIGKILL)
except (OSError, ProcessLookupError):
try:
proc.kill()
except OSError:
pass
threading.Thread(target=_escalate, daemon=True).start()
logger.info("Arrêt de la génération demandé (SIGTERM au pipeline)")
return {"arrêt": "demandé"}
# --- Export multi-dalles (mosaïque image/PDF pour téléphone) ---------------
_export_lock = threading.Lock()
@app.post("/api/export")
def export_tiles(req: ExportRequest):
"""Assemble des dalles adjacentes en image ou PDF (export.py).
Local par nature : les tuiles assemblées sont celles du cache output/
(la webapp légère les possède via le cache à la demande) — aucune
délégation à la machine de traitement. Un seul export à la fois
(CPU/mémoire limités sur Raspberry Pi).
"""
import re as _re
from .export import EXPORT_FORMATS, build_export
fmt = (req.format or "").lower()
if fmt not in EXPORT_FORMATS:
raise HTTPException(400, f"format invalide : {req.format!r} "
f"(attendus : {', '.join(EXPORT_FORMATS)})")
viz = [str(v) for v in (req.viz or []) if v]
if not viz:
raise HTTPException(400, "aucune visualisation demandée")
cells = []
for pair in req.tiles:
if not (isinstance(pair, list) and len(pair) == 2):
raise HTTPException(400, f"tuile invalide : {pair!r} (attendu [col, row])")
cell = (int(pair[0]), int(pair[1]))
if cell not in cells:
cells.append(cell)
if not cells:
raise HTTPException(400, "aucune tuile fournie")
if len(cells) > MAX_EXPORT_TILES:
raise HTTPException(400, f"trop de dalles ({len(cells)}) — max {MAX_EXPORT_TILES}")
if not _export_lock.acquire(blocking=False):
raise HTTPException(409, "un export est déjà en cours — réessayez dans un instant")
try:
out_dir = OUTPUT_DIR / "exports"
result = build_export(OUTPUT_DIR / "visualisations", cells, viz,
float(req.resolution), fmt, out_dir,
max_side=int(req.max_side))
except ValueError as e:
raise HTTPException(400, str(e))
finally:
_export_lock.release()
# Une couche = un fichier image (ou une page du PDF) : la réponse liste
# tous les fichiers générés, chacun avec SA légende.
return {"fichiers": [{
"url": f"/api/export/file/{e['file'].name}",
"nom": e["file"].name,
"taille": e["file"].stat().st_size,
"largeur": e["width"],
"hauteur": e["height"],
"pages": e["pages"],
} for e in result["files"]]}
@app.get("/api/export/file/{name}")
def export_file(name: str):
"""Sert un fichier exporté en téléchargement (pièce jointe)."""
import re as _re
if not _re.fullmatch(r"[A-Za-z0-9_.-]+", name):
raise HTTPException(404, "nom de fichier invalide")
base = (OUTPUT_DIR / "exports").resolve()
path = (base / name).resolve()
if base not in path.parents or not path.is_file():
raise HTTPException(404, f"export introuvable : {name}")
return FileResponse(str(path), filename=name,
headers={"Cache-Control": "no-cache"})
if __name__ == "__main__":
import uvicorn
# HTTPS optionnel : la géolocalisation du navigateur (bouton ⌖, centrage
# GPS depuis un téléphone) n'est disponible qu'en contexte sécurisé.
# Définir LIDAR_SSL_CERTFILE + LIDAR_SSL_KEYFILE (certificat auto-signé)
# pour servir la carte en https://<hôte>:PORT ; le téléphone accepte
# l'avertissement « certificat non fiable » puis le GPS fonctionne.
_ssl_cert = os.environ.get("LIDAR_SSL_CERTFILE")
_ssl_key = os.environ.get("LIDAR_SSL_KEYFILE")
_ssl_kw = ({"ssl_certfile": _ssl_cert, "ssl_keyfile": _ssl_key}
if _ssl_cert and _ssl_key else {})
uvicorn.run(app, host="0.0.0.0", port=PORT, **_ssl_kw)