shuishen
6 hours ago fbb068ec702338d609c1ca6eddbdb9f182d8f211
scripts/serve_workbench_console.py
@@ -5,9 +5,11 @@
import argparse
import base64
import binascii
import hashlib
import json
import os
import re
import shutil
import subprocess
import threading
from datetime import UTC, datetime
@@ -15,24 +17,65 @@
from http.server import SimpleHTTPRequestHandler, ThreadingHTTPServer
from pathlib import Path, PurePosixPath
from typing import Any
from urllib.parse import unquote, urlsplit
from urllib.parse import parse_qs, unquote, urlsplit
from uuid import uuid4
DEFAULT_HOST = "127.0.0.1"
DEFAULT_PORT = 6173
MAX_REQUEST_BYTES = 128 * 1024 * 1024
MAX_FILE_BYTES = 96 * 1024 * 1024
# Uploads are sent as Base64 JSON. Keep the request limit above two 1 GiB
# files after encoding while retaining a per-file bound for local experiments.
MAX_REQUEST_BYTES = 3072 * 1024 * 1024
MAX_FILE_BYTES = 1024 * 1024 * 1024
MAX_IMAGES_PER_RUN = 12
MAX_SEGMENTATION_IMAGES_PER_RUN = 6
MAX_MEASUREMENT_RASTERS_PER_RUN = 4
MAX_POINTCLOUDS_PER_RUN = 2
MAX_PHOTO_RECONSTRUCTION_IMAGES_PER_RUN = 30
PHOTO_RECONSTRUCTION_TIMEOUT = 7200
RISK_RULE_REQUIRED_FILES = {"observations", "zones", "rules"}
MAX_ANOMALY_IMAGES_PER_ROLE = 6
CHANGE_THRESHOLD_DEFAULT = 0.5
CHANGE_THRESHOLD_MIN = 0.01
CHANGE_THRESHOLD_MAX = 0.99
CHANGE_MAX_DIMENSION_DEFAULT = 1024
CHANGE_MAX_DIMENSION_AUTO = 0
CHANGE_MAX_DIMENSION_MIN = 512
CHANGE_MAX_DIMENSION_MAX = 4096
CHANGE_PROCESSING_MODE_DEFAULT = "auto"
CHANGE_PROCESSING_MODES = {"auto", "image", "geotiff"}
SCAN_DEFAULT_THRESHOLDS = [0.3, 0.4, 0.5]
SCAN_DEFAULT_AREAS = [64, 256, 686]
SCAN_MAX_THRESHOLDS = 6
SCAN_MAX_AREAS = 6
SCAN_MAX_COMBINATIONS = 24
SCAN_JOB_TIMEOUT = 1800
ALLOWED_PATH_PREFIXES = (
    "apps/workbench-console",
    "shared/outputs",
    "shared/data/raw/00-change-detection",
    "shared/data/raw/01-object-detection",
    "shared/data/raw/02-semantic-mapping",
    "shared/data/raw/09-anomaly-detection",
    "shared/data/raw/05-3d-pointcloud",
)
SAFE_FILE_NAME = re.compile(r"[^A-Za-z0-9._-]+")
SAFE_FILE_NAME = re.compile(r"[^\w.-]+", re.UNICODE)
SAFE_UPLOAD_ID = re.compile(r"^[0-9a-f]{32}$")
SAFE_SCAN_ID = re.compile(r"^[A-Za-z0-9._-]{1,100}$")
SAFE_SCAN_RESULT_ID = re.compile(r"^threshold-\d+(?:\.\d+)?_area-\d+$")
RUN_LOCK = threading.Lock()
SCAN_JOBS: dict[str, dict[str, Any]] = {}
SCAN_JOBS_LOCK = threading.Lock()
ANOMALY_JOB_LOCK = threading.Lock()
ANOMALY_JOBS: dict[str, dict[str, Any]] = {}
PHOTO_RECONSTRUCTION_JOBS: dict[str, dict[str, Any]] = {}
PHOTO_RECONSTRUCTION_JOBS_LOCK = threading.Lock()
POINTCLOUD_TRAINING_JOBS: dict[str, dict[str, Any]] = {}
POINTCLOUD_TRAINING_JOBS_LOCK = threading.Lock()
POINTCLOUD_INFERENCE_JOBS: dict[str, dict[str, Any]] = {}
POINTCLOUD_INFERENCE_JOBS_LOCK = threading.Lock()
MAX_ANNOTATION_LABELS = 400_000
POINTCLOUD_CLASS_CODES = {1, 2, 5, 6, 15, 16}
class ApiError(ValueError):
@@ -77,6 +120,14 @@
    except (OSError, json.JSONDecodeError):
        return {}
    return payload if isinstance(payload, dict) else {}
def file_sha256(path: Path) -> str:
    digest = hashlib.sha256()
    with path.open("rb") as stream:
        for chunk in iter(lambda: stream.read(8 * 1024 * 1024), b""):
            digest.update(chunk)
    return digest.hexdigest()
def trajectory_runs(root: Path) -> list[dict[str, Any]]:
@@ -131,6 +182,123 @@
    return sorted(records, key=lambda item: (item["createdAt"], item["id"]), reverse=True)
def change_runs(root: Path) -> list[dict[str, Any]]:
    output_root = root / "shared" / "outputs" / "00-change-detection"
    records: list[dict[str, Any]] = []
    for metadata_path in output_root.rglob("run_metadata.json"):
        artifact = metadata_path.parent
        metadata = load_json(metadata_path)
        artifacts = metadata.get("artifacts")
        if metadata.get("capability") != "00-change-detection" or metadata.get("schema_version") != 1 or metadata.get("kind") == "parameter-scan-inference" or not isinstance(artifacts, dict):
            continue
        if not (artifact / str(artifacts.get("overlay") or "")).is_file() or not (artifact / str(artifacts.get("vector") or "")).is_file():
            continue
        raw_root_value = str(metadata.get("raw_input_dir") or "shared/data/raw/00-change-detection/validation-20260817")
        raw_root = root / Path(raw_root_value)
        input_files = metadata.get("input_files")
        if not isinstance(input_files, list) or len(input_files) != 2:
            continue
        before_value = str(metadata.get("raw_before") or (Path(raw_root_value) / str(input_files[0])).as_posix())
        after_value = str(metadata.get("raw_after") or (Path(raw_root_value) / str(input_files[1])).as_posix())
        try:
            before_path = (root / before_value).resolve()
            after_path = (root / after_value).resolve()
            allowed_raw = (root / "shared" / "data" / "raw" / "00-change-detection").resolve()
            before_path.relative_to(allowed_raw)
            after_path.relative_to(allowed_raw)
        except ValueError:
            continue
        if not before_path.is_file() or not after_path.is_file():
            continue
        registered_before_name = str(artifacts.get("before_processed_preview") or "")
        registered_after_name = str(artifacts.get("after_registered_preview") or "")
        registered_before_path = artifact / registered_before_name if registered_before_name else None
        registered_after_path = artifact / registered_after_name if registered_after_name else None
        run_id = artifact.name
        record = {
                "id": run_id,
                "label": run_id,
                "note": "ChangeStar CPU 变化栅格与 GeoAI 像素坐标图斑;结果需人工复核。",
                "artifactRoot": relative_path(root, artifact),
                "beforeImage": relative_path(root, before_path),
                "afterImage": relative_path(root, after_path),
                "createdAt": str(metadata.get("created_at") or ""),
            }
        if registered_before_path and registered_after_path and registered_before_path.is_file() and registered_after_path.is_file():
            record["registeredBeforeImage"] = relative_path(root, registered_before_path)
            record["registeredAfterImage"] = relative_path(root, registered_after_path)
            # Existing console consumers use afterImage as the vector-overlay
            # base. Point it at the registered grid so polygons and highlights
            # share the same pixel coordinates as the model output.
            record["rawAfterImage"] = record["afterImage"]
            record["afterImage"] = record["registeredAfterImage"]
        records.append(record)
    return sorted(records, key=lambda item: (item["createdAt"], item["id"]), reverse=True)
def change_parameter_scans(root: Path) -> list[dict[str, Any]]:
    """Discover read-only parameter scans produced from an existing change run."""
    output_root = root / "shared" / "outputs" / "00-change-detection"
    records: list[dict[str, Any]] = []
    for summary_path in output_root.rglob("scan_summary.json"):
        scan_root = summary_path.parent
        summary = load_json(summary_path)
        results: list[dict[str, Any]] = []
        for item in summary.get("results", []):
            if not isinstance(item, dict) or not isinstance(item.get("directory"), str):
                continue
            directory = scan_root / item["directory"]
            overlay = directory / "overlay_preview.jpg"
            mask = directory / "change_mask.tif"
            regions = directory / "regions.json"
            if not overlay.is_file() or not mask.is_file() or not regions.is_file():
                continue
            results.append(
                {
                    "id": item["directory"],
                    "label": f"T={float(item.get('threshold', 0.5)):.2f} / 面积={int(item.get('minimum_area_pixels', 0))} px",
                    "threshold": item.get("threshold"),
                    "minimumAreaPixels": item.get("minimum_area_pixels"),
                    "cleanedComponents": item.get("cleaned_components"),
                    "changedPixels": item.get("changed_pixels"),
                    "changedPixelRatio": item.get("changed_pixel_ratio"),
                    "vectorFeatureCount": item.get("vector_feature_count"),
                    "fullVectorFeatureCount": (load_json(directory / "full_result.json").get("full_vector_feature_count") if (directory / "full_result.json").is_file() else None),
                    "rectangleFeatureCount": (
                        load_json(directory / "full_result.json").get("rectangle_vector_feature_count")
                        if (directory / "full_result.json").is_file() and load_json(directory / "full_result.json").get("rectangle_vector_feature_count") is not None
                        else len(load_json(directory / "changes_rectangles.geojson").get("features", [])) if (directory / "changes_rectangles.geojson").is_file() else None
                    ),
                    "overlay": relative_path(root, overlay),
                    "mask": relative_path(root, mask),
                    "regions": relative_path(root, regions),
                    "vector": (relative_path(root, directory / "changes.geojson") if (directory / "changes.geojson").is_file() else None),
                    "rectangleVector": (relative_path(root, directory / "changes_rectangles.geojson") if (directory / "changes_rectangles.geojson").is_file() else None),
                    "rectangleVectorWgs84": (relative_path(root, directory / "changes_rectangles_wgs84.geojson") if (directory / "changes_rectangles_wgs84.geojson").is_file() else None),
                }
            )
        # A failed job may have been repaired or materialized later. Keep it
        # discoverable whenever at least one complete candidate exists; only
        # hide scans that still have no usable result.
        if not results:
            continue
        scan_metadata = load_json(scan_root / "scan_metadata.json")
        contact_sheet = scan_root / "parameter_scan_contact_sheet.jpg"
        records.append(
            {
                "id": scan_root.name,
                "label": f"低成本参数扫描 · {scan_root.name}",
                "note": f"复用已有变化概率结果,不重新运行 ChangeStar;源运行:{summary.get('source_run', '未知')}",
                "artifactRoot": relative_path(root, scan_root),
                "sourceRun": summary.get("source_run"),
                "userSubmitted": bool(scan_metadata),
                "contactSheet": relative_path(root, contact_sheet) if contact_sheet.is_file() else None,
                "results": results,
            }
        )
    return sorted(records, key=lambda item: item["id"], reverse=True)
def semantic_runs(root: Path) -> list[dict[str, Any]]:
    output_root = root / "shared" / "outputs" / "02-semantic-mapping"
    records: list[dict[str, Any]] = []
@@ -154,6 +322,420 @@
    return sorted(records, key=lambda item: (item["createdAt"], item["id"]), reverse=True)
def measurement_runs(root: Path) -> list[dict[str, Any]]:
    output_root = root / "shared" / "outputs" / "04-spatial-measurement"
    records: list[dict[str, Any]] = []
    for metadata_path in output_root.rglob("run_metadata.json"):
        artifact = metadata_path.parent
        metadata = load_json(metadata_path)
        if metadata.get("capability") != "04-spatial-measurement" or not isinstance(metadata.get("images"), list):
            continue
        run_id = artifact.name
        records.append(
            {
                "id": run_id,
                "label": run_id,
                "note": "GeoAI 栅格转矢量后进行对象计数、面积和周长测量。",
                "artifactRoot": relative_path(root, artifact),
                "createdAt": str(metadata.get("created_at") or ""),
            }
        )
    return sorted(records, key=lambda item: (item["createdAt"], item["id"]), reverse=True)
def pointcloud_runs(root: Path) -> list[dict[str, Any]]:
    output_root = root / "shared" / "outputs" / "05-3d-pointcloud"
    records: list[dict[str, Any]] = []
    for metadata_path in output_root.rglob("run_metadata.json"):
        artifact = metadata_path.parent
        metadata = load_json(metadata_path)
        dense_photo_reconstruction = metadata.get("dense_photo_reconstruction")
        if metadata.get("capability") == "05-3d-pointcloud" and isinstance(dense_photo_reconstruction, dict):
            textured_model = dense_photo_reconstruction.get("textured_model_file")
            dense_point_cloud = dense_photo_reconstruction.get("dense_point_cloud_file")
            mesh = dense_photo_reconstruction.get("mesh_file")
            if all(isinstance(item, str) and (artifact / item).is_file() for item in (textured_model, dense_point_cloud, mesh)):
                run_id = artifact.name
                records.append(
                    {
                        "id": run_id,
                        "label": str(metadata.get("display_name") or run_id),
                        "note": "CPU 稠密 MVS:显示经过深度融合、网格化和纹理化的局部模型;不是测绘级坐标、DSM、正射图或语义识别结论。",
                        "artifactRoot": relative_path(root, artifact),
                        "createdAt": str(metadata.get("created_at") or ""),
                    }
                )
                continue
        photo_reconstruction = metadata.get("photo_reconstruction")
        if metadata.get("capability") == "05-3d-pointcloud" and isinstance(photo_reconstruction, dict):
            preview = photo_reconstruction.get("preview_file")
            point_cloud = photo_reconstruction.get("point_cloud_file")
            if isinstance(preview, str) and isinstance(point_cloud, str) and (artifact / preview).is_file() and (artifact / point_cloud).is_file():
                run_id = artifact.name
                records.append(
                    {
                        "id": run_id,
                        "label": str(metadata.get("display_name") or run_id),
                        "note": "CPU 稀疏 SfM:显示可复核点云、相机位姿与误差;不是稠密重建、DSM、语义识别或测绘精度结论。",
                        "artifactRoot": relative_path(root, artifact),
                        "createdAt": str(metadata.get("created_at") or ""),
                    }
                )
                continue
        point_clouds = metadata.get("point_clouds")
        if metadata.get("capability") != "05-3d-pointcloud" or not isinstance(point_clouds, list) or not point_clouds:
            continue
        if any(not isinstance(item, dict) or not (artifact / str(item.get("preview_file") or "")).is_file() or not (artifact / str(item.get("vector_file") or "")).is_file() for item in point_clouds):
            continue
        run_id = artifact.name
        records.append(
            {
                "id": run_id,
                "label": str(metadata.get("display_name") or run_id),
                "note": "CPU 语义规则基线:地面、植被、构筑物以及电线/杆塔候选,需要人工复核;不提供测绘精度或资产台账结论。" if any(isinstance(item, dict) and item.get("semantic_summary_file") for item in point_clouds) else "CPU 几何基线:地面/高出地物分离、DSM、近似网格和 GeoAI 足迹;不提供语义类别或测绘精度结论。",
                "artifactRoot": relative_path(root, artifact),
                "createdAt": str(metadata.get("created_at") or ""),
            }
        )
    return sorted(records, key=lambda item: (item["createdAt"], item["id"]), reverse=True)
def pointcloud_annotation_sources(root: Path) -> list[dict[str, Any]]:
    """Expose only generated, fixed preview PLYs suitable for manual labels."""
    sources: list[dict[str, Any]] = []
    for case in pointcloud_runs(root):
        artifact = root / str(case["artifactRoot"])
        metadata = load_json(artifact / "run_metadata.json")
        for cloud in metadata.get("point_clouds", []):
            if not isinstance(cloud, dict):
                continue
            name = cloud.get("semantic_annotation_source_point_cloud")
            if not isinstance(name, str) or Path(name).name != name:
                continue
            path = artifact / name
            if not path.is_file() or path.suffix.lower() != ".ply":
                continue
            sources.append({
                "id": f"{case['id']}:{name}", "runId": case["id"], "label": f"{case['label']} / {cloud.get('file', name)}",
                "artifactRoot": case["artifactRoot"], "file": name, "url": f"/{case['artifactRoot']}/{name}",
                "sha256": file_sha256(path), "pointCount": int(cloud.get("semantic_preview_points") or 0),
                "sourceKind": str(cloud.get("semantic_annotation_source_kind") or "generated point-cloud preview"),
            })
    return sources
def pointcloud_annotations(root: Path) -> list[dict[str, Any]]:
    output = root / "shared" / "outputs" / "05-3d-pointcloud" / "annotations"
    records: list[dict[str, Any]] = []
    for path in output.glob("*/annotation.json"):
        data = load_json(path)
        if data.get("schema_version") != 1 or not isinstance(data.get("id"), str):
            continue
        records.append({"id": data["id"], "sourceId": data.get("source_id"), "createdAt": data.get("created_at"), "labelCount": len(data.get("labels", [])), "classCounts": data.get("class_counts", {}), "path": relative_path(root, path)})
    return sorted(records, key=lambda item: (str(item["createdAt"]), str(item["id"])), reverse=True)
def pointcloud_training_job(job_id: str) -> dict[str, Any] | None:
    with POINTCLOUD_TRAINING_JOBS_LOCK:
        value = POINTCLOUD_TRAINING_JOBS.get(job_id)
        return dict(value) if value else None
def execute_pointcloud_training_job(root: Path, job_id: str, annotation: Path, output: Path, device: str) -> None:
    with POINTCLOUD_TRAINING_JOBS_LOCK:
        POINTCLOUD_TRAINING_JOBS[job_id].update({"status": "running", "stage": "training", "startedAt": datetime.now(UTC).isoformat()})
    python = root / ".venvs" / "05-3d-pointcloud" / "Scripts" / "python.exe"
    command = [str(python), str(root / "capabilities" / "05-3d-pointcloud" / "train_pointcloud_semantic_model.py"), "--annotation", str(annotation), "--output", str(output), "--device", device]
    try:
        with RUN_LOCK:
            completed = subprocess.run(command, cwd=root, capture_output=True, text=True, timeout=14_400, check=False)
        if completed.returncode:
            message = (completed.stderr or completed.stdout or "Unknown training error.").strip().splitlines()[-1]
            raise ApiError(message[:600])
        metrics = output / "metrics.json"
        model = output / "model.pt"
        preview = output / "predicted-semantic-preview.ply"
        if not all(path.is_file() for path in (metrics, model, preview)):
            raise ApiError("Training finished without model, metrics, and predicted preview artifacts.")
        with POINTCLOUD_TRAINING_JOBS_LOCK:
            POINTCLOUD_TRAINING_JOBS[job_id].update({"status": "complete", "stage": "complete", "completedAt": datetime.now(UTC).isoformat(), "artifactRoot": relative_path(root, output), "metrics": relative_path(root, metrics), "model": relative_path(root, model), "preview": relative_path(root, preview)})
    except Exception as exc:
        with POINTCLOUD_TRAINING_JOBS_LOCK:
            POINTCLOUD_TRAINING_JOBS[job_id].update({"status": "failed", "stage": "failed", "completedAt": datetime.now(UTC).isoformat(), "error": str(exc)[:700]})
def pointcloud_semantic_models(root: Path) -> list[dict[str, Any]]:
    """Expose only complete locally trained models, never arbitrary model paths."""
    output_root = root / "shared" / "outputs" / "05-3d-pointcloud" / "training-runs"
    records: list[dict[str, Any]] = []
    for model_path in output_root.glob("*/model.pt"):
        metrics_path = model_path.with_name("metrics.json")
        metrics = load_json(metrics_path)
        classes = metrics.get("classes")
        if metrics.get("capability") != "05-3d-pointcloud" or metrics.get("classification") != "B" or not isinstance(classes, dict):
            continue
        class_codes = sorted(str(code) for code in classes if str(code).isdigit())
        if len(class_codes) < 2:
            continue
        test = metrics.get("test") if isinstance(metrics.get("test"), dict) else {}
        report = test.get("report") if isinstance(test.get("report"), dict) else {}
        summary: dict[str, float] = {}
        for code in class_codes:
            definition = classes.get(code)
            key = definition.get("key") if isinstance(definition, dict) else None
            score = report.get(key) if isinstance(key, str) else None
            if isinstance(score, dict) and isinstance(score.get("f1-score"), (int, float)):
                summary[key] = round(float(score["f1-score"]), 3)
        records.append({"id": model_path.parent.name, "label": model_path.parent.name, "artifactRoot": relative_path(root, model_path.parent), "model": relative_path(root, model_path), "metrics": relative_path(root, metrics_path), "createdAt": str(metrics.get("created_at") or ""), "classes": classes, "testF1": summary})
    return sorted(records, key=lambda item: (item["createdAt"], item["id"]), reverse=True)
def pointcloud_inference_job(job_id: str) -> dict[str, Any] | None:
    with POINTCLOUD_INFERENCE_JOBS_LOCK:
        value = POINTCLOUD_INFERENCE_JOBS.get(job_id)
        return dict(value) if value else None
def execute_pointcloud_inference_job(root: Path, job_id: str, model: Path, source: Path, output: Path) -> None:
    with POINTCLOUD_INFERENCE_JOBS_LOCK:
        POINTCLOUD_INFERENCE_JOBS[job_id].update({"status": "running", "stage": "inference", "startedAt": datetime.now(UTC).isoformat()})
    python = root / ".venvs" / "05-3d-pointcloud" / "Scripts" / "python.exe"
    command = [str(python), str(root / "capabilities" / "05-3d-pointcloud" / "apply_pointcloud_semantic_model.py"), "--model", str(model), "--input", str(source), "--output", str(output), "--device", "cpu"]
    try:
        with RUN_LOCK:
            completed = subprocess.run(command, cwd=root, capture_output=True, text=True, timeout=14_400, check=False)
        if completed.returncode:
            message = (completed.stderr or completed.stdout or "Unknown model inference error.").strip().splitlines()[-1]
            raise ApiError(message[:600])
        metadata = output / "run_metadata.json"
        preview = output / "predicted-semantic-preview.ply"
        classified_las = output / "predicted-semantic-classified.las"
        counts = output / "class-counts.csv"
        summary = output / "prediction-summary.json"
        if not all(path.is_file() for path in (metadata, preview, classified_las, counts, summary)):
            raise ApiError("Model inference finished without all expected prediction artifacts.")
        with POINTCLOUD_INFERENCE_JOBS_LOCK:
            POINTCLOUD_INFERENCE_JOBS[job_id].update({"status": "complete", "stage": "complete", "completedAt": datetime.now(UTC).isoformat(), "artifactRoot": relative_path(root, output), "metadata": relative_path(root, metadata), "preview": relative_path(root, preview), "classifiedLas": relative_path(root, classified_las), "classCounts": relative_path(root, counts), "summary": relative_path(root, summary)})
    except Exception as exc:
        with POINTCLOUD_INFERENCE_JOBS_LOCK:
            POINTCLOUD_INFERENCE_JOBS[job_id].update({"status": "failed", "stage": "failed", "completedAt": datetime.now(UTC).isoformat(), "error": str(exc)[:700]})
def risk_rule_runs(root: Path) -> list[dict[str, Any]]:
    output_root = root / "shared" / "outputs" / "07-risk-rule-engine"
    records: list[dict[str, Any]] = []
    for metadata_path in output_root.rglob("run_metadata.json"):
        artifact = metadata_path.parent
        metadata = load_json(metadata_path)
        artifacts = metadata.get("artifacts")
        if metadata.get("capability") != "07-risk-rule-engine" or not isinstance(artifacts, dict):
            continue
        required = ("risk_raster", "risk_preview", "risk_vector", "risk_scores_csv", "summary")
        if any(not isinstance(artifacts.get(key), str) or not (artifact / artifacts[key]).is_file() for key in required):
            continue
        run_id = artifact.name
        records.append(
            {
                "id": run_id,
                "label": str(metadata.get("display_name") or run_id),
                "note": "可审计空间规则评分,仅供人工复核,不构成事件或处置结论。",
                "artifactRoot": relative_path(root, artifact),
                "createdAt": str(metadata.get("created_at") or ""),
            }
        )
    return sorted(records, key=lambda item: (item["createdAt"], item["id"]), reverse=True)
def anomaly_runs(root: Path) -> list[dict[str, Any]]:
    output_root = root / "shared" / "outputs" / "09-anomaly-detection"
    allowed_raw = (root / "shared" / "data" / "raw" / "09-anomaly-detection").resolve()
    records: list[dict[str, Any]] = []
    for metadata_path in output_root.rglob("run_metadata.json"):
        artifact = metadata_path.parent
        metadata = load_json(metadata_path)
        images = metadata.get("images")
        if metadata.get("capability") != "09-anomaly-detection" or not isinstance(images, list):
            continue
        raw_input_value = str(metadata.get("raw_input_dir") or "")
        raw_reference_value = str(metadata.get("raw_reference_dir") or "")
        if not raw_input_value or not raw_reference_value:
            continue
        try:
            raw_input = (root / raw_input_value).resolve()
            raw_reference = (root / raw_reference_value).resolve()
            raw_input.relative_to(allowed_raw)
            raw_reference.relative_to(allowed_raw)
        except ValueError:
            continue
        if not raw_input.is_dir() or not raw_reference.is_dir():
            continue
        if any(not (artifact / str(item.get("overlay_file") or "")).is_file() for item in images if isinstance(item, dict)):
            continue
        run_id = artifact.name
        records.append(
            {
                "id": run_id,
                "label": str(metadata.get("display_name") or run_id),
                "note": str(metadata.get("case_note") or "规则基线与 Isolation Forest 的视觉离群候选,只供人工复核。"),
                "artifactRoot": relative_path(root, artifact),
                "inputRoot": relative_path(root, raw_input),
                "referenceRoot": relative_path(root, raw_reference),
                "createdAt": str(metadata.get("created_at") or ""),
            }
        )
    return sorted(records, key=lambda item: (item["createdAt"], item["id"]), reverse=True)
def anomaly_job(job_id: str) -> dict[str, Any] | None:
    with ANOMALY_JOB_LOCK:
        value = ANOMALY_JOBS.get(job_id)
        return dict(value) if value else None
def photo_reconstruction_job(job_id: str) -> dict[str, Any] | None:
    with PHOTO_RECONSTRUCTION_JOBS_LOCK:
        value = PHOTO_RECONSTRUCTION_JOBS.get(job_id)
        return dict(value) if value else None
def run_background_command(command: list[str], root: Path, timeout: int) -> None:
    completed = subprocess.run(command, cwd=root, capture_output=True, text=True, timeout=timeout, check=False)
    if completed.returncode:
        message = (completed.stderr or completed.stdout or "Unknown script error.").strip().splitlines()[-1]
        raise ApiError(f"Processing failed: {message[:600]}")
def execute_photo_reconstruction_job(
    root: Path,
    job_id: str,
    run_id: str,
    raw_root: Path,
    processed_root: Path,
    sparse_output: Path,
    output: Path,
    source_sha256: dict[str, str],
    source_bytes: dict[str, int],
    use_position_priors: bool,
) -> None:
    python = root / ".venvs" / "05-3d-pointcloud" / "Scripts" / "python.exe"
    sparse_command = [
        str(python), str(root / "capabilities" / "05-3d-pointcloud" / "run_photo_reconstruction.py"),
        "--input", str(processed_root), "--output", str(sparse_output),
        "--max-image-size", "2000", "--max-features", "18000",
        "--camera-model", "OPENCV",
    ]
    if use_position_priors:
        sparse_command.extend(["--matching-mode", "spatial", "--matching-neighbors", "4", "--use-position-priors", "--prior-position-loss-scale-m", "0.05"])
    else:
        sparse_command.extend(["--matching-mode", "exhaustive"])
    dense_command = [
        str(python), str(root / "capabilities" / "05-3d-pointcloud" / "run_cpu_dense_reconstruction.py"),
        "--input", str(processed_root), "--sparse-model", str(sparse_output / "sparse_model" / "0"),
        "--output", str(output),
        "--openmvs-bin", str(root / "shared" / "tools" / "openmvs-2.4.0" / "vc17" / "x64" / "Release"),
        "--threads", "12", "--max-resolution", "2400", "--dense-resolution-level", "0",
        "--dense-number-views", "8", "--dense-number-views-fuse", "2", "--target-faces", "800000",
    ]
    try:
        with PHOTO_RECONSTRUCTION_JOBS_LOCK:
            PHOTO_RECONSTRUCTION_JOBS[job_id].update({"status": "running", "stage": "sparse_sfm", "startedAt": datetime.now(UTC).isoformat()})
        with RUN_LOCK:
            run_background_command(sparse_command, root, 1800)
            sparse_metadata = load_json(sparse_output / "run_metadata.json")
            if not (sparse_output / "sparse_model" / "0").is_dir():
                raise ApiError("Sparse photo reconstruction finished without the expected COLMAP model.")
            with PHOTO_RECONSTRUCTION_JOBS_LOCK:
                PHOTO_RECONSTRUCTION_JOBS[job_id].update({"stage": "dense_mvs"})
            run_background_command(dense_command, root, PHOTO_RECONSTRUCTION_TIMEOUT)
        metadata_path = output / "run_metadata.json"
        if not metadata_path.is_file():
            raise ApiError("CPU dense reconstruction finished without the expected result metadata.")
        metadata = load_json(metadata_path)
        metadata["photo_reconstruction"] = sparse_metadata.get("photo_reconstruction", {})
        metadata["input_dir"] = relative_path(root, processed_root)
        metadata["raw_input_dir"] = relative_path(root, raw_root)
        metadata["source_sha256"] = source_sha256
        metadata["source_bytes"] = source_bytes
        metadata["console_photo_reconstruction"] = {"use_position_priors": use_position_priors, "matching_mode": "spatial" if use_position_priors else "exhaustive"}
        metadata["display_name"] = f"用户照片 CPU 稠密重建({len(source_sha256)} 图)"
        metadata["case_note"] = "用户上传的同架次 JPG/JPEG 照片经 CPU SfM/MVS 重建;需要人工检查几何与纹理质量,不是测绘级成果。"
        metadata_path.write_text(json.dumps(metadata, ensure_ascii=False, indent=2), encoding="utf-8")
        definition = next(item for item in pointcloud_runs(root) if item["id"] == run_id)
        with PHOTO_RECONSTRUCTION_JOBS_LOCK:
            PHOTO_RECONSTRUCTION_JOBS[job_id].update({"status": "complete", "stage": "complete", "run": definition, "finishedAt": datetime.now(UTC).isoformat()})
    except Exception as exc:  # pragma: no cover - background boundary
        with PHOTO_RECONSTRUCTION_JOBS_LOCK:
            PHOTO_RECONSTRUCTION_JOBS[job_id].update({"status": "failed", "stage": "failed", "error": str(exc), "finishedAt": datetime.now(UTC).isoformat()})
def validate_anomaly_parameters(payload: dict[str, Any]) -> tuple[int, int, float, int]:
    tile_size = payload.get("tileSize", 256)
    stride = payload.get("stride", 128)
    threshold_quantile = payload.get("thresholdQuantile", 0.995)
    random_state = payload.get("randomState", 42)
    if isinstance(tile_size, bool) or not isinstance(tile_size, int) or not 128 <= tile_size <= 1024:
        raise ApiError("Tile size must be an integer between 128 and 1024.")
    if isinstance(stride, bool) or not isinstance(stride, int) or not 32 <= stride <= tile_size:
        raise ApiError("Stride must be an integer between 32 and tile size.")
    if isinstance(threshold_quantile, bool) or not isinstance(threshold_quantile, (int, float)) or not 0.9 <= float(threshold_quantile) <= 0.9999:
        raise ApiError("Threshold quantile must be between 0.9 and 0.9999.")
    if isinstance(random_state, bool) or not isinstance(random_state, int) or not 0 <= random_state <= 2_147_483_647:
        raise ApiError("Random state must be a non-negative integer.")
    return tile_size, stride, float(threshold_quantile), random_state
def execute_anomaly_job(
    root: Path,
    job_id: str,
    run_id: str,
    raw_reference: Path,
    raw_input: Path,
    processed_reference: Path,
    processed_input: Path,
    output: Path,
    tile_size: int,
    stride: int,
    threshold_quantile: float,
    random_state: int,
) -> None:
    with ANOMALY_JOB_LOCK:
        ANOMALY_JOBS[job_id]["status"] = "running"
    python = root / ".venvs" / "09-anomaly-detection" / "Scripts" / "python.exe"
    command = [
        str(python),
        str(root / "capabilities" / "09-anomaly-detection" / "run_anomaly_detection.py"),
        "--reference", str(processed_reference),
        "--input", str(processed_input),
        "--output", str(output),
        "--tile-size", str(tile_size),
        "--stride", str(stride),
        "--threshold-quantile", f"{threshold_quantile:.6f}",
        "--random-state", str(random_state),
        "--spatial-mode", "auto",
    ]
    try:
        with RUN_LOCK:
            completed = subprocess.run(command, cwd=root, capture_output=True, text=True, timeout=1800, check=False)
        if completed.returncode:
            message = (completed.stderr or completed.stdout or "Unknown script error.").strip().splitlines()[-1]
            raise ApiError(f"Processing failed: {message[:600]}")
        metadata_path = output / "run_metadata.json"
        if not metadata_path.is_file():
            raise ApiError("Anomaly-detection script finished without the expected result metadata.")
        metadata = load_json(metadata_path)
        metadata["raw_input_dir"] = relative_path(root, raw_input)
        metadata["raw_reference_dir"] = relative_path(root, raw_reference)
        metadata["processed_input_dir"] = relative_path(root, processed_input)
        metadata["processed_reference_dir"] = relative_path(root, processed_reference)
        metadata_path.write_text(json.dumps(metadata, ensure_ascii=False, indent=2), encoding="utf-8")
        definition = next(item for item in anomaly_runs(root) if item["id"] == run_id)
        with ANOMALY_JOB_LOCK:
            ANOMALY_JOBS[job_id].update({"status": "complete", "run": definition, "finishedAt": datetime.now(UTC).isoformat()})
    except Exception as exc:  # pragma: no cover - background boundary
        with ANOMALY_JOB_LOCK:
            ANOMALY_JOBS[job_id].update({"status": "failed", "error": str(exc), "finishedAt": datetime.now(UTC).isoformat()})
def semantic_tasks(root: Path) -> list[dict[str, Any]]:
    catalog = load_json(root / "capabilities" / "02-semantic-mapping" / "configs" / "task-catalog.json")
    tasks = catalog.get("tasks")
@@ -173,6 +755,21 @@
    def do_GET(self) -> None:  # noqa: N802 - inherited standard-library method name
        path = urlsplit(self.path).path
        if path == "/api/change-detection/runs":
            self.send_json(HTTPStatus.OK, {"runs": change_runs(self.root)})
            return
        if path == "/api/change-detection/scans":
            self.send_json(HTTPStatus.OK, {"scans": change_parameter_scans(self.root)})
            return
        if path.startswith("/api/change-detection/scan-jobs/"):
            job_id = path.rstrip("/").rsplit("/", 1)[-1]
            with SCAN_JOBS_LOCK:
                job = dict(SCAN_JOBS.get(job_id, {}))
            if not job:
                self.send_json(HTTPStatus.NOT_FOUND, {"error": "Unknown change-detection scan job."})
            else:
                self.send_json(HTTPStatus.OK, {"job": job})
            return
        if path == "/api/trajectory/runs":
            self.send_json(HTTPStatus.OK, {"runs": trajectory_runs(self.root)})
            return
@@ -185,6 +782,47 @@
        if path == "/api/semantic-mapping/tasks":
            self.send_json(HTTPStatus.OK, {"tasks": semantic_tasks(self.root)})
            return
        if path == "/api/spatial-measurement/runs":
            self.send_json(HTTPStatus.OK, {"runs": measurement_runs(self.root)})
            return
        if path == "/api/3d-pointcloud/runs":
            self.send_json(HTTPStatus.OK, {"runs": pointcloud_runs(self.root)})
            return
        if path == "/api/3d-pointcloud/annotation-sources":
            self.send_json(HTTPStatus.OK, {"sources": pointcloud_annotation_sources(self.root)})
            return
        if path == "/api/3d-pointcloud/annotations":
            self.send_json(HTTPStatus.OK, {"annotations": pointcloud_annotations(self.root)})
            return
        if path == "/api/3d-pointcloud/semantic-models":
            self.send_json(HTTPStatus.OK, {"models": pointcloud_semantic_models(self.root)})
            return
        if path.startswith("/api/3d-pointcloud/model-inference-jobs/"):
            job_id = path.rstrip("/").rsplit("/", 1)[-1]
            job = pointcloud_inference_job(job_id)
            self.send_json(HTTPStatus.OK if job else HTTPStatus.NOT_FOUND, {"job": job} if job else {"error": "Unknown point-cloud model inference job."})
            return
        if path.startswith("/api/3d-pointcloud/training-jobs/"):
            job_id = path.rstrip("/").rsplit("/", 1)[-1]
            job = pointcloud_training_job(job_id)
            self.send_json(HTTPStatus.OK if job else HTTPStatus.NOT_FOUND, {"job": job} if job else {"error": "Unknown point-cloud training job."})
            return
        if path.startswith("/api/3d-pointcloud/photo-reconstruction-jobs/"):
            job_id = path.rstrip("/").rsplit("/", 1)[-1]
            job = photo_reconstruction_job(job_id)
            self.send_json(HTTPStatus.OK if job else HTTPStatus.NOT_FOUND, {"job": job} if job else {"error": "Unknown photo-reconstruction job."})
            return
        if path == "/api/risk-rule-engine/runs":
            self.send_json(HTTPStatus.OK, {"runs": risk_rule_runs(self.root)})
            return
        if path == "/api/anomaly-detection/runs":
            self.send_json(HTTPStatus.OK, {"runs": anomaly_runs(self.root)})
            return
        if path.startswith("/api/anomaly-detection/jobs/"):
            job_id = path.rstrip("/").rsplit("/", 1)[-1]
            job = anomaly_job(job_id)
            self.send_json(HTTPStatus.OK if job else HTTPStatus.NOT_FOUND, {"job": job} if job else {"error": "Unknown anomaly-detection job."})
            return
        if path == "/":
            self.send_response(HTTPStatus.FOUND)
            self.send_header("Location", "/apps/workbench-console/")
@@ -196,6 +834,16 @@
        path = urlsplit(self.path).path
        try:
            payload = self.read_json_body()
            if path == "/api/change-detection/runs":
                self.send_json(HTTPStatus.CREATED, {"run": self.create_change_run(payload)})
                return
            if path == "/api/change-detection/scans":
                self.send_json(HTTPStatus.ACCEPTED, {"job": self.create_change_scan(payload)})
                return
            if path.startswith("/api/change-detection/scans/") and path.endswith("/promote"):
                scan_id = path.split("/")[-2]
                self.send_json(HTTPStatus.CREATED, {"run": self.promote_change_scan(scan_id, payload)})
                return
            if path == "/api/trajectory/runs":
                self.send_json(HTTPStatus.CREATED, {"run": self.create_trajectory_run(payload)})
                return
@@ -204,6 +852,30 @@
                return
            if path == "/api/semantic-mapping/runs":
                self.send_json(HTTPStatus.CREATED, {"run": self.create_semantic_run(payload)})
                return
            if path == "/api/spatial-measurement/runs":
                self.send_json(HTTPStatus.CREATED, {"run": self.create_measurement_run(payload)})
                return
            if path == "/api/3d-pointcloud/runs":
                self.send_json(HTTPStatus.CREATED, {"run": self.create_pointcloud_run(payload)})
                return
            if path == "/api/3d-pointcloud/annotations":
                self.send_json(HTTPStatus.CREATED, {"annotation": self.create_pointcloud_annotation(payload)})
                return
            if path == "/api/3d-pointcloud/training-runs":
                self.send_json(HTTPStatus.ACCEPTED, {"job": self.create_pointcloud_training_run(payload)})
                return
            if path == "/api/3d-pointcloud/model-inference-runs":
                self.send_json(HTTPStatus.ACCEPTED, {"job": self.create_pointcloud_model_inference_run(payload)})
                return
            if path == "/api/3d-pointcloud/photo-reconstruction-runs":
                self.send_json(HTTPStatus.ACCEPTED, {"job": self.create_photo_reconstruction_run(payload)})
                return
            if path == "/api/risk-rule-engine/runs":
                self.send_json(HTTPStatus.CREATED, {"run": self.create_risk_rule_run(payload)})
                return
            if path == "/api/anomaly-detection/runs":
                self.send_json(HTTPStatus.ACCEPTED, {"job": self.create_anomaly_run(payload)})
                return
            self.send_json(HTTPStatus.NOT_FOUND, {"error": "Unknown local API endpoint."})
        except ApiError as exc:
@@ -214,9 +886,58 @@
            self.log_error("local run failed: %s", exc)
            self.send_json(HTTPStatus.INTERNAL_SERVER_ERROR, {"error": "Local run failed. Check the console terminal for details."})
    def do_DELETE(self) -> None:  # noqa: N802 - annotation revisions are explicitly user-removable
        path = urlsplit(self.path).path
        prefix = "/api/3d-pointcloud/annotations/"
        if not path.startswith(prefix):
            self.send_json(HTTPStatus.NOT_FOUND, {"error": "Unknown local API endpoint."})
            return
        try:
            annotation_id = path[len(prefix):]
            if not annotation_id or "/" in annotation_id or SAFE_FILE_NAME.search(annotation_id) or len(annotation_id) > 120:
                raise ApiError("Invalid annotation id.")
            record = next((item for item in pointcloud_annotations(self.root) if item["id"] == annotation_id), None)
            if not record:
                raise ApiError("The selected annotation revision is unavailable.")
            annotation_root = (self.root / "shared" / "outputs" / "05-3d-pointcloud" / "annotations").resolve()
            location = (annotation_root / annotation_id).resolve()
            location.relative_to(annotation_root)
            if not (location / "annotation.json").is_file():
                raise ApiError("The selected annotation revision is incomplete.")
            shutil.rmtree(location)
            self.send_json(HTTPStatus.OK, {"deletedId": annotation_id})
        except ApiError as exc:
            self.send_json(HTTPStatus.BAD_REQUEST, {"error": str(exc)})
        except ValueError:
            self.send_json(HTTPStatus.BAD_REQUEST, {"error": "Invalid annotation location."})
        except Exception as exc:  # pragma: no cover - defensive server boundary
            self.log_error("annotation deletion failed: %s", exc)
            self.send_json(HTTPStatus.INTERNAL_SERVER_ERROR, {"error": "Annotation deletion failed. Check the console terminal for details."})
    def do_PUT(self) -> None:  # noqa: N802 - binary upload endpoint
        path = urlsplit(self.path).path
        if not path.startswith("/api/change-detection/uploads/") and not path.startswith("/api/anomaly-detection/uploads/") and not path.startswith("/api/3d-pointcloud/photo-uploads/") and not path.startswith("/api/3d-pointcloud/pointcloud-uploads/"):
            self.send_json(HTTPStatus.NOT_FOUND, {"error": "Unknown local API endpoint."})
            return
        try:
            if path.startswith("/api/anomaly-detection/uploads/"):
                result = self.receive_anomaly_upload(path)
            elif path.startswith("/api/3d-pointcloud/photo-uploads/"):
                result = self.receive_photo_reconstruction_upload(path)
            elif path.startswith("/api/3d-pointcloud/pointcloud-uploads/"):
                result = self.receive_pointcloud_upload(path)
            else:
                result = self.receive_change_upload(path)
            self.send_json(HTTPStatus.CREATED, result)
        except ApiError as exc:
            self.send_json(HTTPStatus.BAD_REQUEST, {"error": str(exc)})
        except Exception as exc:  # pragma: no cover - defensive server boundary
            self.log_error("binary upload failed: %s", exc)
            self.send_json(HTTPStatus.INTERNAL_SERVER_ERROR, {"error": "Binary upload failed. Check the console terminal for details."})
    def do_OPTIONS(self) -> None:  # noqa: N802
        self.send_response(HTTPStatus.NO_CONTENT)
        self.send_header("Allow", "GET, POST, OPTIONS")
        self.send_header("Allow", "GET, POST, PUT, DELETE, OPTIONS")
        self.end_headers()
    def read_json_body(self) -> dict[str, Any]:
@@ -241,6 +962,102 @@
        if completed.returncode:
            message = (completed.stderr or completed.stdout or "Unknown script error.").strip().splitlines()[-1]
            raise ApiError(f"Processing failed: {message[:600]}")
    def receive_change_upload(self, path: str) -> dict[str, Any]:
        return self.receive_binary_upload(path, "00-change-detection", {"before", "after"}, {".jpg", ".jpeg", ".png", ".tif", ".tiff"}, "change-detection")
    def receive_anomaly_upload(self, path: str) -> dict[str, Any]:
        return self.receive_binary_upload(path, "09-anomaly-detection", {"reference", "input"}, {".jpg", ".jpeg", ".png", ".tif", ".tiff"}, "anomaly-detection")
    def receive_photo_reconstruction_upload(self, path: str) -> dict[str, Any]:
        return self.receive_binary_upload(path, "05-3d-pointcloud", {"photo"}, {".jpg", ".jpeg"}, "photo reconstruction")
    def receive_pointcloud_upload(self, path: str) -> dict[str, Any]:
        return self.receive_binary_upload(path, "05-3d-pointcloud", {"pointcloud"}, {".ply", ".pcd", ".xyz", ".xyzn", ".xyzrgb", ".las", ".laz"}, "point-cloud")
    def receive_binary_upload(
        self,
        path: str,
        capability: str,
        allowed_roles: set[str],
        suffixes: set[str],
        label: str,
    ) -> dict[str, Any]:
        upload_id = path.rstrip("/").rsplit("/", 1)[-1]
        if not SAFE_UPLOAD_ID.fullmatch(upload_id):
            raise ApiError(f"Invalid {label} upload id.")
        query = parse_qs(urlsplit(self.path).query)
        role = query.get("role", [""])[0]
        if role not in allowed_roles:
            raise ApiError(f"Invalid {label} upload role.")
        encoded_name = self.headers.get("X-Upload-Name", "")
        if len(encoded_name) > 2048:
            raise ApiError("Encoded upload name is too long.")
        try:
            name = unquote(encoded_name, encoding="utf-8", errors="strict")
        except UnicodeError as exc:
            raise ApiError("Upload name is not valid UTF-8 percent encoding.") from exc
        safe_name = safe_file_name(name, suffixes)
        content_length = self.headers.get("Content-Length")
        if content_length is None or not content_length.isdigit():
            raise ApiError("Binary upload requires a Content-Length header.")
        size = int(content_length)
        if size <= 0 or size > MAX_FILE_BYTES:
            raise ApiError(f"Uploaded file must be between 1 byte and {MAX_FILE_BYTES // (1024 * 1024)} MB: {safe_name}.")
        staging = self.root / "shared" / "data" / "raw" / capability / "uploads" / upload_id
        staging.mkdir(parents=True, exist_ok=False)
        part = staging / f"{role}.part"
        target = staging / f"{role}{Path(safe_name).suffix.lower()}"
        remaining = size
        digest = hashlib.sha256()
        try:
            with part.open("wb") as stream:
                while remaining:
                    chunk = self.rfile.read(min(8 * 1024 * 1024, remaining))
                    if not chunk:
                        raise ApiError("Binary upload ended before Content-Length was reached.")
                    stream.write(chunk)
                    digest.update(chunk)
                    remaining -= len(chunk)
            part.replace(target)
            (staging / f"{role}.json").write_text(json.dumps({"role": role, "name": safe_name, "size": size, "sha256": digest.hexdigest()}), encoding="utf-8")
        except Exception:
            part.unlink(missing_ok=True)
            target.unlink(missing_ok=True)
            raise
        return {"uploadId": upload_id, "role": role, "name": safe_name, "size": size, "sha256": digest.hexdigest()}
    def resolve_change_upload(self, payload: Any, role: str) -> tuple[str, Path]:
        return self.resolve_binary_upload(payload, role, "00-change-detection", "change-detection")
    def resolve_anomaly_upload(self, payload: Any, role: str) -> tuple[str, Path, str]:
        name, path = self.resolve_binary_upload(payload, role, "09-anomaly-detection", "anomaly-detection")
        manifest = load_json(path.parent / f"{role}.json")
        return name, path, str(manifest.get("sha256") or "")
    def resolve_photo_reconstruction_upload(self, payload: Any) -> tuple[str, Path, str]:
        name, path = self.resolve_binary_upload(payload, "photo", "05-3d-pointcloud", "photo reconstruction")
        manifest = load_json(path.parent / "photo.json")
        return name, path, str(manifest.get("sha256") or "")
    def resolve_pointcloud_upload(self, payload: Any) -> tuple[str, Path, str]:
        name, path = self.resolve_binary_upload(payload, "pointcloud", "05-3d-pointcloud", "point-cloud")
        manifest = load_json(path.parent / "pointcloud.json")
        return name, path, str(manifest.get("sha256") or "")
    def resolve_binary_upload(self, payload: Any, role: str, capability: str, label: str) -> tuple[str, Path]:
        if not isinstance(payload, dict) or not isinstance(payload.get("uploadId"), str):
            raise ApiError(f"{label} uploads must include a {role} uploadId.")
        upload_id = payload["uploadId"]
        if not SAFE_UPLOAD_ID.fullmatch(upload_id):
            raise ApiError(f"Invalid {label} upload id.")
        staging = self.root / "shared" / "data" / "raw" / capability / "uploads" / upload_id
        manifest = load_json(staging / f"{role}.json")
        name = str(manifest.get("name") or "")
        path = staging / f"{role}{Path(name).suffix.lower()}"
        if manifest.get("role") != role or not name or not path.is_file():
            raise ApiError(f"The staged {role} upload is unavailable or incomplete.")
        return name, path
    def create_trajectory_run(self, payload: dict[str, Any]) -> dict[str, Any]:
        files = payload.get("files")
@@ -299,6 +1116,327 @@
            raise ApiError("Detection script finished without the expected result metadata.")
        return next(item for item in detection_runs(self.root) if item["id"] == run_id)
    def create_change_run(self, payload: dict[str, Any]) -> dict[str, Any]:
        files = payload.get("files")
        uploads = payload.get("uploads")
        staged: dict[str, tuple[str, Path]] = {}
        if isinstance(uploads, dict):
            staged["before"] = self.resolve_change_upload(uploads.get("before"), "before")
            staged["after"] = self.resolve_change_upload(uploads.get("after"), "after")
        elif isinstance(files, dict):
            decoded_before = decode_upload(files.get("before"), {".jpg", ".jpeg", ".png", ".tif", ".tiff"})
            decoded_after = decode_upload(files.get("after"), {".jpg", ".jpeg", ".png", ".tif", ".tiff"})
        else:
            raise ApiError("Change-detection request must contain before and after files or uploads.")
        threshold_value = payload.get("threshold", CHANGE_THRESHOLD_DEFAULT)
        if isinstance(threshold_value, bool) or not isinstance(threshold_value, (int, float)):
            raise ApiError("Change-detection threshold must be a number between 0.01 and 0.99.")
        threshold = float(threshold_value)
        if not CHANGE_THRESHOLD_MIN <= threshold <= CHANGE_THRESHOLD_MAX:
            raise ApiError("Change-detection threshold must be between 0.01 and 0.99.")
        processing_mode = payload.get("processingMode", CHANGE_PROCESSING_MODE_DEFAULT)
        if not isinstance(processing_mode, str) or processing_mode not in CHANGE_PROCESSING_MODES:
            raise ApiError("Change-detection processing mode must be auto, image, or geotiff.")
        max_dimension_value = payload.get("maxDimension", CHANGE_MAX_DIMENSION_AUTO)
        if isinstance(max_dimension_value, bool) or not isinstance(max_dimension_value, int):
            raise ApiError("Change-detection resolution must be an integer: 0 or between 512 and 4096.")
        max_dimension = int(max_dimension_value)
        if max_dimension != CHANGE_MAX_DIMENSION_AUTO and not CHANGE_MAX_DIMENSION_MIN <= max_dimension <= CHANGE_MAX_DIMENSION_MAX:
            raise ApiError("Change-detection resolution must be 0 or between 512 and 4096.")
        if staged:
            before_name, after_name = staged["before"][0], staged["after"][0]
        else:
            before_name, after_name = decoded_before[0], decoded_after[0]
        run_id = make_run_id("change")
        raw_root = self.root / "shared" / "data" / "raw" / "00-change-detection" / "runs" / run_id
        before_path = raw_root / "before" / before_name
        after_path = raw_root / "after" / after_name
        before_path.parent.mkdir(parents=True, exist_ok=False)
        after_path.parent.mkdir(parents=True, exist_ok=False)
        if staged:
            shutil.copyfile(staged["before"][1], before_path)
            shutil.copyfile(staged["after"][1], after_path)
        else:
            before_path.write_bytes(decoded_before[1])
            after_path.write_bytes(decoded_after[1])
        processed_root = self.root / "shared" / "data" / "processed" / "00-change-detection" / run_id
        output = self.root / "shared" / "outputs" / "00-change-detection" / "runs" / run_id
        python = self.root / ".venvs" / "00-change-detection" / "Scripts" / "python.exe"
        if not python.is_file():
            raise ApiError("Change-detection virtual environment is unavailable. Run the capability setup first.")
        with RUN_LOCK:
            self.run_command(
                [
                    str(python),
                    str(self.root / "capabilities" / "00-change-detection" / "run_change_detection.py"),
                    "--before", str(before_path),
                    "--after", str(after_path),
                    "--threshold", f"{threshold:.4f}",
                    "--max-dimension", str(max_dimension),
                    "--processing-mode", processing_mode,
                    "--processed-output", str(processed_root),
                    "--output", str(output),
                ],
                1200,
            )
        metadata_path = output / "run_metadata.json"
        if not metadata_path.is_file():
            raise ApiError("Change-detection script finished without the expected result metadata.")
        metadata = load_json(metadata_path)
        metadata["raw_input_dir"] = relative_path(self.root, raw_root)
        metadata["processed_input_dir"] = relative_path(self.root, processed_root)
        metadata["raw_before"] = relative_path(self.root, before_path)
        metadata["raw_after"] = relative_path(self.root, after_path)
        metadata_path.write_text(json.dumps(metadata, ensure_ascii=False, indent=2), encoding="utf-8")
        return next(item for item in change_runs(self.root) if item["id"] == run_id)
    def _scan_parameters(self, payload: dict[str, Any]) -> tuple[list[float], list[int]]:
        raw_thresholds = payload.get("thresholds", SCAN_DEFAULT_THRESHOLDS)
        raw_areas = payload.get("minimumAreas", SCAN_DEFAULT_AREAS)
        if not isinstance(raw_thresholds, list) or not raw_thresholds or len(raw_thresholds) > SCAN_MAX_THRESHOLDS:
            raise ApiError(f"Parameter scan thresholds must contain 1-{SCAN_MAX_THRESHOLDS} values.")
        if not isinstance(raw_areas, list) or not raw_areas or len(raw_areas) > SCAN_MAX_AREAS:
            raise ApiError(f"Parameter scan minimum areas must contain 1-{SCAN_MAX_AREAS} values.")
        thresholds: list[float] = []
        for value in raw_thresholds:
            if isinstance(value, bool) or not isinstance(value, (int, float)):
                raise ApiError("Each scan threshold must be a number between 0.01 and 0.99.")
            number = round(float(value), 4)
            if not CHANGE_THRESHOLD_MIN <= number <= CHANGE_THRESHOLD_MAX:
                raise ApiError("Each scan threshold must be between 0.01 and 0.99.")
            if number not in thresholds:
                thresholds.append(number)
        areas: list[int] = []
        for value in raw_areas:
            if isinstance(value, bool) or not isinstance(value, int) or not 16 <= value <= 200000:
                raise ApiError("Each scan minimum area must be an integer between 16 and 200000 pixels.")
            if value not in areas:
                areas.append(value)
        if len(thresholds) * len(areas) > SCAN_MAX_COMBINATIONS:
            raise ApiError(f"A parameter scan accepts at most {SCAN_MAX_COMBINATIONS} combinations.")
        return thresholds, areas
    def _scan_root(self, scan_id: str) -> Path:
        if not SAFE_SCAN_ID.fullmatch(scan_id):
            raise ApiError("Invalid parameter-scan id.")
        scan_root = self.root / "shared" / "outputs" / "00-change-detection" / "parameter-scans" / scan_id
        if not scan_root.is_dir() or not (scan_root / "scan_summary.json").is_file():
            raise ApiError("The parameter-scan result is unavailable.")
        return scan_root
    def promote_change_scan(self, scan_id: str, payload: dict[str, Any]) -> dict[str, Any]:
        scan_root = self._scan_root(scan_id)
        result_id = payload.get("resultId")
        if not isinstance(result_id, str) or not SAFE_SCAN_RESULT_ID.fullmatch(result_id):
            raise ApiError("A valid parameter-scan resultId is required.")
        result_dir = scan_root / result_id
        summary = load_json(result_dir / "summary.json")
        full_result = load_json(result_dir / "full_result.json")
        inference_run_id = str(load_json(scan_root / "scan_summary.json").get("source_run") or "")
        inference_output = self.root / "shared" / "outputs" / "00-change-detection" / "runs" / inference_run_id
        inference_metadata = load_json(inference_output / "run_metadata.json")
        if not result_dir.is_dir() or not (result_dir / "changes.geojson").is_file() or not inference_metadata:
            raise ApiError("The selected scan result is incomplete and cannot be promoted.")
        run_id = make_run_id("change")
        output = self.root / "shared" / "outputs" / "00-change-detection" / "runs" / run_id
        output.mkdir(parents=True, exist_ok=False)
        for source_name, destination_name in (
            ("change_probability.tif", "change_probability.tif"),
            ("change_mask.tif", "change_mask.tif"),
            ("generic_difference_mask.tif", "generic_difference_mask.tif"),
            ("change_model_overlay.jpg", "change_model_overlay.jpg"),
            ("before_processed_preview.jpg", "before_processed_preview.jpg"),
            ("after_registered_preview.jpg", "after_registered_preview.jpg"),
            ("changes.geojson", "changes.geojson"),
            ("changes_rectangles.geojson", "changes_rectangles.geojson"),
            ("changes_rectangles_wgs84.geojson", "changes_rectangles_wgs84.geojson"),
        ):
            source = result_dir / source_name if source_name.startswith("change_mask") or source_name.startswith("changes") else inference_output / source_name
            if source.is_file():
                shutil.copyfile(source, output / destination_name)
        overlay_source = result_dir / "overlay_preview.jpg"
        if (inference_output / "change_overlay.jpg").is_file() and float(summary.get("threshold", 0.5)) == float(inference_metadata.get("thresholds", {}).get("change_probability", 0.5)):
            overlay_source = inference_output / "change_overlay.jpg"
        if not overlay_source.is_file():
            raise ApiError("The selected scan preview is unavailable.")
        shutil.copyfile(overlay_source, output / "change_overlay.jpg")
        metadata = dict(inference_metadata)
        scan_raw_root = self.root / "shared" / "data" / "raw" / "00-change-detection" / "runs" / scan_id
        before_raw = scan_raw_root / "before" / str(inference_metadata.get("input_files", ["before.tif", "after.tif"])[0])
        after_raw = scan_raw_root / "after" / str(inference_metadata.get("input_files", ["before.tif", "after.tif"])[1])
        rectangle_count = int(full_result.get("rectangle_vector_feature_count") or 0)
        if rectangle_count == 0 and (result_dir / "changes_rectangles.geojson").is_file():
            rectangle_count = len(load_json(result_dir / "changes_rectangles.geojson").get("features", []))
        metadata.update(
            {
                "kind": "formal-change-run",
                "created_at": datetime.now(UTC).isoformat(),
                "thresholds": {"change_probability": float(summary.get("threshold", 0.5)), "minimum_component_pixels": int(summary.get("minimum_area_pixels", 16))},
                "raw_changed_pixels": int(summary.get("raw_changed_pixels", 0)),
                "changed_pixels": int(summary.get("changed_pixels", 0)),
                "changed_pixel_ratio": float(summary.get("changed_pixel_ratio", 0)),
                "vector_feature_count": int(full_result.get("full_vector_feature_count", summary.get("vector_feature_count", 0))),
                "rectangle_feature_count": rectangle_count,
                "promoted_from_scan": scan_id,
                "promoted_result": result_id,
                "raw_input_dir": relative_path(self.root, scan_raw_root),
                "raw_before": relative_path(self.root, before_raw),
                "raw_after": relative_path(self.root, after_raw),
                "processed_input_dir": relative_path(self.root, self.root / "shared" / "data" / "processed" / "00-change-detection" / scan_id),
                "generic_difference": inference_metadata.get("generic_difference"),
                "artifacts": {
                    "probability_raster": "change_probability.tif",
                    "raw_mask_raster": "change_mask.tif",
                    "mask_raster": "change_mask.tif",
                    "generic_difference_mask": "generic_difference_mask.tif" if (output / "generic_difference_mask.tif").is_file() else None,
                    "overlay": "change_overlay.jpg",
                    "model_overlay": "change_model_overlay.jpg" if (output / "change_model_overlay.jpg").is_file() else None,
                    "before_processed_preview": "before_processed_preview.jpg" if (output / "before_processed_preview.jpg").is_file() else None,
                    "after_registered_preview": "after_registered_preview.jpg" if (output / "after_registered_preview.jpg").is_file() else None,
                    "vector": "changes.geojson",
                    "rectangle_vector": "changes_rectangles.geojson",
                    "rectangle_vector_wgs84": "changes_rectangles_wgs84.geojson" if (output / "changes_rectangles_wgs84.geojson").is_file() else None,
                    "features": "full_result.json",
                },
            }
        )
        (output / "full_result.json").write_text(json.dumps(full_result, ensure_ascii=False, indent=2), encoding="utf-8")
        (output / "run_metadata.json").write_text(json.dumps(metadata, ensure_ascii=False, indent=2), encoding="utf-8")
        return next(item for item in change_runs(self.root) if item["id"] == run_id)
    def create_change_scan(self, payload: dict[str, Any]) -> dict[str, Any]:
        uploads = payload.get("uploads")
        if not isinstance(uploads, dict):
            raise ApiError("Parameter scan must contain staged before and after uploads.")
        staged = {
            "before": self.resolve_change_upload(uploads.get("before"), "before"),
            "after": self.resolve_change_upload(uploads.get("after"), "after"),
        }
        thresholds, areas = self._scan_parameters(payload)
        processing_mode = payload.get("processingMode", CHANGE_PROCESSING_MODE_DEFAULT)
        if not isinstance(processing_mode, str) or processing_mode not in CHANGE_PROCESSING_MODES:
            raise ApiError("Change-detection processing mode must be auto, image, or geotiff.")
        max_dimension_value = payload.get("maxDimension", CHANGE_MAX_DIMENSION_AUTO)
        if isinstance(max_dimension_value, bool) or not isinstance(max_dimension_value, int):
            raise ApiError("Change-detection resolution must be an integer: 0 or between 512 and 4096.")
        max_dimension = int(max_dimension_value)
        if max_dimension != CHANGE_MAX_DIMENSION_AUTO and not CHANGE_MAX_DIMENSION_MIN <= max_dimension <= CHANGE_MAX_DIMENSION_MAX:
            raise ApiError("Change-detection resolution must be 0 or between 512 and 4096.")
        run_id = make_run_id("scan")
        raw_root = self.root / "shared" / "data" / "raw" / "00-change-detection" / "runs" / run_id
        before_path = raw_root / "before" / staged["before"][0]
        after_path = raw_root / "after" / staged["after"][0]
        before_path.parent.mkdir(parents=True, exist_ok=False)
        after_path.parent.mkdir(parents=True, exist_ok=False)
        shutil.copyfile(staged["before"][1], before_path)
        shutil.copyfile(staged["after"][1], after_path)
        processed_root = self.root / "shared" / "data" / "processed" / "00-change-detection" / run_id
        inference_output = self.root / "shared" / "outputs" / "00-change-detection" / "runs" / run_id
        scan_output = self.root / "shared" / "outputs" / "00-change-detection" / "parameter-scans" / run_id
        with SCAN_JOBS_LOCK:
            SCAN_JOBS[run_id] = {
                "id": run_id,
                "status": "queued",
                "createdAt": datetime.now(UTC).isoformat(),
                "thresholds": thresholds,
                "minimumAreas": areas,
                "processingMode": processing_mode,
                "maxDimension": max_dimension,
            }
        thread = threading.Thread(
            target=self._run_change_scan,
            args=(run_id, before_path, after_path, processed_root, inference_output, scan_output, thresholds, areas, processing_mode, max_dimension),
            daemon=True,
            name=f"change-scan-{run_id}",
        )
        thread.start()
        return dict(SCAN_JOBS[run_id])
    def _update_scan_job(self, job_id: str, **values: Any) -> None:
        with SCAN_JOBS_LOCK:
            if job_id in SCAN_JOBS:
                SCAN_JOBS[job_id].update(values)
    def _run_change_scan(
        self,
        run_id: str,
        before_path: Path,
        after_path: Path,
        processed_root: Path,
        inference_output: Path,
        scan_output: Path,
        thresholds: list[float],
        areas: list[int],
        processing_mode: str,
        max_dimension: int,
    ) -> None:
        python = self.root / ".venvs" / "00-change-detection" / "Scripts" / "python.exe"
        try:
            if not python.is_file():
                raise ApiError("Change-detection virtual environment is unavailable. Run the capability setup first.")
            self._update_scan_job(run_id, status="running", phase="inference")
            with RUN_LOCK:
                self.run_command(
                    [
                        str(python),
                        str(self.root / "capabilities" / "00-change-detection" / "run_change_detection.py"),
                        "--before", str(before_path),
                        "--after", str(after_path),
                        "--threshold", "0.5000",
                        "--max-dimension", str(max_dimension),
                        "--processing-mode", processing_mode,
                        "--processed-output", str(processed_root),
                        "--output", str(inference_output),
                    ],
                    SCAN_JOB_TIMEOUT,
                )
                inference_metadata_path = inference_output / "run_metadata.json"
                inference_metadata = load_json(inference_metadata_path)
                inference_metadata["kind"] = "parameter-scan-inference"
                inference_metadata["scan_job_id"] = run_id
                inference_metadata_path.write_text(json.dumps(inference_metadata, ensure_ascii=False, indent=2), encoding="utf-8")
                self._update_scan_job(run_id, phase="parameter-scan")
                command = [
                    str(python),
                    str(self.root / "capabilities" / "00-change-detection" / "scan_change_detection_parameters.py"),
                    "--run-dir", str(inference_output),
                    "--output", str(scan_output),
                ]
                for threshold in thresholds:
                    command.extend(["--threshold", f"{threshold:.4f}"])
                for area in areas:
                    command.extend(["--minimum-area", str(area)])
                self.run_command(command, SCAN_JOB_TIMEOUT)
                self._update_scan_job(run_id, phase="vectorization")
                vector_command = [
                    str(python),
                    str(self.root / "capabilities" / "00-change-detection" / "materialize_parameter_scan_candidates.py"),
                    "--scan-dir", str(scan_output),
                ]
                for threshold in thresholds:
                    for area in areas:
                        vector_command.extend(["--candidate", f"threshold-{threshold:.2f}_area-{area}"])
                self.run_command(vector_command, SCAN_JOB_TIMEOUT)
            metadata = {
                "capability": "00-change-detection",
                "kind": "parameter-scan",
                "source_run": run_id,
                "created_at": datetime.now(UTC).isoformat(),
                "thresholds": thresholds,
                "minimum_areas": areas,
                "processing_mode": processing_mode,
                "max_dimension": max_dimension,
                "raw_input_dir": relative_path(self.root, before_path.parent.parent),
                "processed_input_dir": relative_path(self.root, processed_root),
            }
            scan_output.mkdir(parents=True, exist_ok=True)
            (scan_output / "scan_metadata.json").write_text(json.dumps(metadata, ensure_ascii=False, indent=2), encoding="utf-8")
            self._update_scan_job(run_id, status="completed", phase="done", scanId=run_id)
        except Exception as exc:  # background errors are returned through polling
            scan_output.mkdir(parents=True, exist_ok=True)
            (scan_output / "scan_failed.json").write_text(json.dumps({"job_id": run_id, "error": str(exc)[:600]}, ensure_ascii=False, indent=2), encoding="utf-8")
            self._update_scan_job(run_id, status="failed", phase="error", error=str(exc)[:600])
    def create_semantic_run(self, payload: dict[str, Any]) -> dict[str, Any]:
        task_id = str(payload.get("taskId") or "color_baseline")
        task = next((item for item in semantic_tasks(self.root) if item["id"] == task_id), None)
@@ -339,6 +1477,360 @@
        metadata_path.write_text(json.dumps(metadata, ensure_ascii=False, indent=2), encoding="utf-8")
        return next(item for item in semantic_runs(self.root) if item["id"] == run_id)
    def create_anomaly_run(self, payload: dict[str, Any]) -> dict[str, Any]:
        tile_size, stride, threshold_quantile, random_state = validate_anomaly_parameters(payload)
        uploads = payload.get("uploads")
        if not isinstance(uploads, dict):
            raise ApiError("Anomaly-detection request must contain reference and input uploads.")
        reference_values = uploads.get("reference")
        input_values = uploads.get("input")
        if not isinstance(reference_values, list) or not reference_values:
            raise ApiError("Select at least one normal reference image.")
        if not isinstance(input_values, list) or not input_values:
            raise ApiError("Select at least one image to inspect.")
        if len(reference_values) > MAX_ANOMALY_IMAGES_PER_ROLE or len(input_values) > MAX_ANOMALY_IMAGES_PER_ROLE:
            raise ApiError(f"An anomaly-detection run accepts at most {MAX_ANOMALY_IMAGES_PER_ROLE} images in each group.")
        references = [self.resolve_anomaly_upload(value, "reference") for value in reference_values]
        inputs = [self.resolve_anomaly_upload(value, "input") for value in input_values]
        if len({name.casefold() for name, _, _ in references}) != len(references):
            raise ApiError("Normal reference image names must be unique within one run.")
        if len({name.casefold() for name, _, _ in inputs}) != len(inputs):
            raise ApiError("Input image names must be unique within one run.")
        python = self.root / ".venvs" / "09-anomaly-detection" / "Scripts" / "python.exe"
        if not python.is_file():
            raise ApiError("Anomaly-detection virtual environment is unavailable. Run the capability setup first.")
        run_id = make_run_id("anomaly")
        job_id = uuid4().hex
        raw_root = self.root / "shared" / "data" / "raw" / "09-anomaly-detection" / "runs" / run_id
        raw_reference = raw_root / "reference"
        raw_input = raw_root / "input"
        processed_root = self.root / "shared" / "data" / "processed" / "09-anomaly-detection" / run_id
        processed_reference = processed_root / "reference"
        processed_input = processed_root / "input"
        output = self.root / "shared" / "outputs" / "09-anomaly-detection" / "runs" / run_id
        for directory in (raw_reference, raw_input, processed_reference, processed_input):
            directory.mkdir(parents=True, exist_ok=False)
        for group, raw_dir, processed_dir in ((references, raw_reference, processed_reference), (inputs, raw_input, processed_input)):
            for name, staged_path, expected_sha256 in group:
                raw_path = raw_dir / name
                processed_path = processed_dir / name
                shutil.copyfile(staged_path, raw_path)
                if expected_sha256 and file_sha256(raw_path) != expected_sha256:
                    raise ApiError(f"Uploaded file checksum changed while staging: {name}.")
                shutil.copyfile(raw_path, processed_path)
        for _, staged_path, _ in references + inputs:
            shutil.rmtree(staged_path.parent)
        created_at = datetime.now(UTC).isoformat()
        job = {"id": job_id, "runId": run_id, "status": "queued", "createdAt": created_at}
        with ANOMALY_JOB_LOCK:
            ANOMALY_JOBS[job_id] = job
        thread = threading.Thread(
            target=execute_anomaly_job,
            args=(self.root, job_id, run_id, raw_reference, raw_input, processed_reference, processed_input, output, tile_size, stride, threshold_quantile, random_state),
            daemon=True,
            name=f"anomaly-{run_id}",
        )
        thread.start()
        return dict(job)
    def create_measurement_run(self, payload: dict[str, Any]) -> dict[str, Any]:
        uploads = payload.get("rasters")
        if not isinstance(uploads, list) or not uploads:
            raise ApiError("Spatial-measurement request must include at least one label raster.")
        if len(uploads) > MAX_MEASUREMENT_RASTERS_PER_RUN:
            raise ApiError(f"A spatial-measurement run accepts at most {MAX_MEASUREMENT_RASTERS_PER_RUN} rasters.")
        decoded = [decode_upload(item, {".png", ".tif", ".tiff"}) for item in uploads]
        if len({name.casefold() for name, _ in decoded}) != len(decoded):
            raise ApiError("Uploaded raster names must be unique within one run.")
        run_id = make_run_id("measurement")
        raw_root = self.root / "shared" / "data" / "raw" / "04-spatial-measurement" / "runs" / run_id
        processed_root = self.root / "shared" / "data" / "processed" / "04-spatial-measurement" / run_id
        raw_root.mkdir(parents=True, exist_ok=False)
        processed_root.mkdir(parents=True, exist_ok=False)
        for name, content in decoded:
            (raw_root / name).write_bytes(content)
            (processed_root / name).write_bytes(content)
        output = self.root / "shared" / "outputs" / "04-spatial-measurement" / "runs" / run_id
        python = self.root / ".venvs" / "04-spatial-measurement" / "Scripts" / "python.exe"
        if not python.is_file():
            raise ApiError("Spatial-measurement virtual environment is unavailable. Run the capability setup first.")
        with RUN_LOCK:
            self.run_command([str(python), str(self.root / "capabilities" / "04-spatial-measurement" / "run_spatial_measurement.py"), "--input", str(processed_root), "--output", str(output)], 900)
        metadata_path = output / "run_metadata.json"
        if not metadata_path.is_file():
            raise ApiError("Spatial-measurement script finished without the expected result metadata.")
        metadata = load_json(metadata_path)
        metadata["input_dir"] = relative_path(self.root, processed_root)
        metadata["raw_input_dir"] = relative_path(self.root, raw_root)
        metadata_path.write_text(json.dumps(metadata, ensure_ascii=False, indent=2), encoding="utf-8")
        return next(item for item in measurement_runs(self.root) if item["id"] == run_id)
    def create_pointcloud_run(self, payload: dict[str, Any]) -> dict[str, Any]:
        uploads = payload.get("pointClouds")
        source_dense_run_id = payload.get("sourceDenseRunId")
        if uploads is None and not isinstance(source_dense_run_id, str):
            raise ApiError("3D point-cloud request must include at least one PLY, PCD, XYZ, LAS, or LAZ file.")
        if uploads is not None and (not isinstance(uploads, list) or not uploads):
            raise ApiError("3D point-cloud request must include at least one PLY, PCD, XYZ, LAS, or LAZ file.")
        if isinstance(uploads, list) and len(uploads) > MAX_POINTCLOUDS_PER_RUN:
            raise ApiError(f"A 3D point-cloud run accepts at most {MAX_POINTCLOUDS_PER_RUN} files.")
        suffixes = {".ply", ".pcd", ".xyz", ".xyzn", ".xyzrgb", ".las", ".laz"}
        staged_upload_dirs: list[Path] = []
        if isinstance(source_dense_run_id, str):
            if SAFE_FILE_NAME.search(source_dense_run_id) or len(source_dense_run_id) > 120:
                raise ApiError("Invalid dense point-cloud source run id.")
            source_case = next((item for item in pointcloud_runs(self.root) if item["id"] == source_dense_run_id), None)
            if not source_case:
                raise ApiError("The selected dense point-cloud source is unavailable.")
            source_artifact = self.root / str(source_case["artifactRoot"])
            source_metadata = load_json(source_artifact / "run_metadata.json")
            dense = source_metadata.get("dense_photo_reconstruction")
            source_file = dense.get("dense_point_cloud_file") if isinstance(dense, dict) else None
            source_path = source_artifact / str(source_file or "")
            if not isinstance(source_file, str) or source_path.suffix.lower() != ".ply" or not source_path.is_file():
                raise ApiError("The selected run has no available dense PLY output.")
            decoded = [(f"{source_dense_run_id}-dense.ply", source_path, file_sha256(source_path))]
        elif all(isinstance(item, dict) and isinstance(item.get("content"), str) for item in uploads):
            decoded = [(name, content, "") for name, content in (decode_upload(item, suffixes) for item in uploads)]
        else:
            decoded = [self.resolve_pointcloud_upload(item) for item in uploads]
            staged_upload_dirs = [path.parent for _, path, _ in decoded]
        if len({name.casefold() for name, _, _ in decoded}) != len(decoded):
            raise ApiError("Uploaded point-cloud names must be unique within one run.")
        run_id = make_run_id("pointcloud")
        raw_root = self.root / "shared" / "data" / "raw" / "05-3d-pointcloud" / "runs" / run_id
        processed_root = self.root / "shared" / "data" / "processed" / "05-3d-pointcloud" / run_id
        raw_root.mkdir(parents=True, exist_ok=False)
        processed_root.mkdir(parents=True, exist_ok=False)
        source_sha256: dict[str, str] = {}
        source_bytes: dict[str, int] = {}
        for name, staged_or_content, expected_sha256 in decoded:
            raw_path = raw_root / name
            if isinstance(staged_or_content, bytes):
                raw_path.write_bytes(staged_or_content)
            else:
                shutil.copyfile(staged_or_content, raw_path)
            if expected_sha256 and file_sha256(raw_path) != expected_sha256:
                raise ApiError(f"Uploaded point-cloud checksum changed while staging: {name}.")
            source_sha256[name] = file_sha256(raw_path)
            source_bytes[name] = raw_path.stat().st_size
            shutil.copyfile(raw_path, processed_root / name)
        output = self.root / "shared" / "outputs" / "05-3d-pointcloud" / "runs" / run_id
        python = self.root / ".venvs" / "05-3d-pointcloud" / "Scripts" / "python.exe"
        if not python.is_file():
            raise ApiError("3D point-cloud virtual environment is unavailable. Run the capability setup first.")
        command = [str(python), str(self.root / "capabilities" / "05-3d-pointcloud" / "run_pointcloud_understanding.py"), "--input", str(processed_root), "--output", str(output), "--ground-up-axis", "z"]
        with RUN_LOCK:
            self.run_command(command, 900)
        metadata_path = output / "run_metadata.json"
        if not metadata_path.is_file():
            raise ApiError("3D point-cloud script finished without the expected result metadata.")
        metadata = load_json(metadata_path)
        metadata["input_dir"] = relative_path(self.root, processed_root)
        metadata["raw_input_dir"] = relative_path(self.root, raw_root)
        metadata["source_sha256"] = source_sha256
        metadata["source_bytes"] = source_bytes
        metadata_path.write_text(json.dumps(metadata, ensure_ascii=False, indent=2), encoding="utf-8")
        for staging in staged_upload_dirs:
            shutil.rmtree(staging)
        return next(item for item in pointcloud_runs(self.root) if item["id"] == run_id)
    def create_pointcloud_annotation(self, payload: dict[str, Any]) -> dict[str, Any]:
        source_id = payload.get("sourceId")
        labels = payload.get("labels")
        if not isinstance(source_id, str) or not isinstance(labels, list):
            raise ApiError("Annotation request must include a sourceId and labels array.")
        source = next((item for item in pointcloud_annotation_sources(self.root) if item["id"] == source_id), None)
        if not source:
            raise ApiError("The selected generated annotation source is unavailable.")
        if len(labels) > MAX_ANNOTATION_LABELS:
            raise ApiError(f"An annotation revision accepts at most {MAX_ANNOTATION_LABELS} labelled points.")
        compact: dict[int, int] = {}
        for item in labels:
            if not isinstance(item, list) or len(item) != 2 or not all(isinstance(value, int) for value in item):
                raise ApiError("Each annotation label must be [pointIndex, classCode].")
            index, code = item
            if index < 0 or index >= int(source["pointCount"]) or code not in POINTCLOUD_CLASS_CODES:
                raise ApiError("Annotation contains an out-of-range point index or unsupported class code.")
            compact[index] = code
        if not compact:
            raise ApiError("Save at least one user-confirmed point label.")
        annotation_id = make_run_id("annotation")
        location = self.root / "shared" / "outputs" / "05-3d-pointcloud" / "annotations" / annotation_id
        location.mkdir(parents=True, exist_ok=False)
        source_path = self.root / str(source["artifactRoot"]) / str(source["file"])
        class_counts = {str(code): sum(value == code for value in compact.values()) for code in sorted(POINTCLOUD_CLASS_CODES)}
        document = {
            "schema_version": 1, "id": annotation_id, "created_at": datetime.now(UTC).isoformat(),
            "source_id": source_id, "source_path": str(source_path.resolve()), "source_sha256": source["sha256"],
            "source_run_id": source["runId"], "point_count": int(source["pointCount"]),
            "labels": [[index, code] for index, code in sorted(compact.items())], "class_counts": class_counts,
            "class_schema": {str(code): {"code": code} for code in sorted(POINTCLOUD_CLASS_CODES)},
            "provenance": "human_confirmed_point_labels_only",
        }
        path = location / "annotation.json"
        path.write_text(json.dumps(document, ensure_ascii=False, indent=2), encoding="utf-8")
        return {"id": annotation_id, "sourceId": source_id, "path": relative_path(self.root, path), "labelCount": len(compact), "classCounts": class_counts, "createdAt": document["created_at"]}
    def create_pointcloud_training_run(self, payload: dict[str, Any]) -> dict[str, Any]:
        annotation_id = payload.get("annotationId")
        device = payload.get("device", "auto")
        if not isinstance(annotation_id, str) or SAFE_FILE_NAME.search(annotation_id) or len(annotation_id) > 120:
            raise ApiError("Invalid annotation id.")
        if device not in {"auto", "cpu", "cuda"}:
            raise ApiError("Training device must be auto, cpu, or cuda.")
        annotation = self.root / "shared" / "outputs" / "05-3d-pointcloud" / "annotations" / annotation_id / "annotation.json"
        record = load_json(annotation)
        if record.get("schema_version") != 1:
            raise ApiError("The selected annotation revision is unavailable.")
        python = self.root / ".venvs" / "05-3d-pointcloud" / "Scripts" / "python.exe"
        if not python.is_file():
            raise ApiError("3D point-cloud virtual environment is unavailable. Run the capability setup first.")
        job_id = uuid4().hex
        output = self.root / "shared" / "outputs" / "05-3d-pointcloud" / "training-runs" / make_run_id("semantic-model")
        job = {"id": job_id, "annotationId": annotation_id, "status": "queued", "stage": "queued", "device": device, "createdAt": datetime.now(UTC).isoformat()}
        with POINTCLOUD_TRAINING_JOBS_LOCK:
            POINTCLOUD_TRAINING_JOBS[job_id] = job
        thread = threading.Thread(target=execute_pointcloud_training_job, args=(self.root, job_id, annotation, output, device), daemon=True, name=f"pointcloud-training-{job_id[:8]}")
        thread.start()
        return dict(job)
    def create_pointcloud_model_inference_run(self, payload: dict[str, Any]) -> dict[str, Any]:
        model_id = payload.get("modelId")
        upload = payload.get("pointCloud")
        if not isinstance(model_id, str) or SAFE_FILE_NAME.search(model_id) or len(model_id) > 120:
            raise ApiError("Invalid trained model id.")
        model_record = next((item for item in pointcloud_semantic_models(self.root) if item["id"] == model_id), None)
        if not model_record:
            raise ApiError("The selected trained model is unavailable or incomplete.")
        name, staged_path, expected_sha256 = self.resolve_pointcloud_upload(upload)
        suffixes = {".ply", ".pcd", ".xyz", ".xyzn", ".xyzrgb", ".las", ".laz"}
        if Path(name).suffix.lower() not in suffixes:
            raise ApiError("Model inference requires a PLY, PCD, XYZ, LAS, or LAZ point cloud.")
        python = self.root / ".venvs" / "05-3d-pointcloud" / "Scripts" / "python.exe"
        if not python.is_file():
            raise ApiError("3D point-cloud virtual environment is unavailable. Run the capability setup first.")
        run_id = make_run_id("semantic-inference")
        raw_root = self.root / "shared" / "data" / "raw" / "05-3d-pointcloud" / "model-inference-runs" / run_id
        processed_root = self.root / "shared" / "data" / "processed" / "05-3d-pointcloud" / "model-inference-runs" / run_id
        output = self.root / "shared" / "outputs" / "05-3d-pointcloud" / "model-inference-runs" / run_id
        raw_root.mkdir(parents=True, exist_ok=False)
        processed_root.mkdir(parents=True, exist_ok=False)
        raw_path = raw_root / name
        shutil.copyfile(staged_path, raw_path)
        actual_sha256 = file_sha256(raw_path)
        if expected_sha256 and actual_sha256 != expected_sha256:
            raise ApiError("Uploaded point-cloud checksum changed while staging.")
        processed_path = processed_root / name
        shutil.copyfile(raw_path, processed_path)
        # Only remove the staging copy after its immutable raw copy was verified.
        shutil.rmtree(staged_path.parent)
        model_path = self.root / str(model_record["model"])
        training_root = (self.root / "shared" / "outputs" / "05-3d-pointcloud" / "training-runs").resolve()
        try:
            model_path.resolve().relative_to(training_root)
        except ValueError as exc:
            raise ApiError("Selected model is outside the allowed training output directory.") from exc
        job_id = uuid4().hex
        job = {"id": job_id, "runId": run_id, "modelId": model_id, "inputName": name, "status": "queued", "stage": "queued", "device": "cpu", "createdAt": datetime.now(UTC).isoformat(), "sourceSha256": actual_sha256, "rawInput": relative_path(self.root, raw_path), "processedInput": relative_path(self.root, processed_path)}
        with POINTCLOUD_INFERENCE_JOBS_LOCK:
            POINTCLOUD_INFERENCE_JOBS[job_id] = job
        thread = threading.Thread(target=execute_pointcloud_inference_job, args=(self.root, job_id, model_path, processed_path, output), daemon=True, name=f"pointcloud-inference-{job_id[:8]}")
        thread.start()
        return dict(job)
    def create_photo_reconstruction_run(self, payload: dict[str, Any]) -> dict[str, Any]:
        uploads = payload.get("photos")
        if not isinstance(uploads, list) or len(uploads) < 3:
            raise ApiError("Photo reconstruction needs at least three JPG/JPEG photos from one coherent flight or camera sequence.")
        if len(uploads) > MAX_PHOTO_RECONSTRUCTION_IMAGES_PER_RUN:
            raise ApiError(f"A photo reconstruction run accepts at most {MAX_PHOTO_RECONSTRUCTION_IMAGES_PER_RUN} photos.")
        use_position_priors = payload.get("usePositionPriors", False)
        if not isinstance(use_position_priors, bool):
            raise ApiError("Photo reconstruction usePositionPriors must be true or false.")
        photos = [self.resolve_photo_reconstruction_upload(value) for value in uploads]
        if len({name.casefold() for name, _, _ in photos}) != len(photos):
            raise ApiError("Uploaded photo names must be unique within one run.")
        python = self.root / ".venvs" / "05-3d-pointcloud" / "Scripts" / "python.exe"
        openmvs = self.root / "shared" / "tools" / "openmvs-2.4.0" / "vc17" / "x64" / "Release"
        if not python.is_file() or not (openmvs / "DensifyPointCloud.exe").is_file():
            raise ApiError("Photo-reconstruction CPU environment is unavailable. Run the capability setup first.")
        run_id = make_run_id("photo-reconstruction")
        job_id = uuid4().hex
        raw_root = self.root / "shared" / "data" / "raw" / "05-3d-pointcloud" / "runs" / run_id
        processed_root = self.root / "shared" / "data" / "processed" / "05-3d-pointcloud" / run_id
        sparse_output = processed_root / "sparse_sfm"
        output = self.root / "shared" / "outputs" / "05-3d-pointcloud" / "runs" / run_id
        raw_root.mkdir(parents=True, exist_ok=False)
        processed_root.mkdir(parents=True, exist_ok=False)
        source_sha256: dict[str, str] = {}
        source_bytes: dict[str, int] = {}
        for name, staged_path, expected_sha256 in photos:
            raw_path = raw_root / name
            shutil.copyfile(staged_path, raw_path)
            actual_sha256 = file_sha256(raw_path)
            if expected_sha256 and actual_sha256 != expected_sha256:
                raise ApiError(f"Uploaded file checksum changed while staging: {name}.")
            shutil.copyfile(raw_path, processed_root / name)
            source_sha256[name] = actual_sha256
            source_bytes[name] = raw_path.stat().st_size
        for _, staged_path, _ in photos:
            shutil.rmtree(staged_path.parent)
        job = {"id": job_id, "runId": run_id, "status": "queued", "stage": "queued", "createdAt": datetime.now(UTC).isoformat(), "inputImages": len(photos), "usePositionPriors": use_position_priors}
        with PHOTO_RECONSTRUCTION_JOBS_LOCK:
            PHOTO_RECONSTRUCTION_JOBS[job_id] = job
        thread = threading.Thread(
            target=execute_photo_reconstruction_job,
            args=(self.root, job_id, run_id, raw_root, processed_root, sparse_output, output, source_sha256, source_bytes, use_position_priors),
            daemon=True,
            name=f"photo-reconstruction-{run_id}",
        )
        thread.start()
        return dict(job)
    def create_risk_rule_run(self, payload: dict[str, Any]) -> dict[str, Any]:
        files = payload.get("files")
        if not isinstance(files, dict) or set(files) != RISK_RULE_REQUIRED_FILES:
            raise ApiError("Risk-rule request must contain observations, zones and rules files.")
        observations_name, observations_bytes = decode_upload(files["observations"], {".geojson"})
        zones_name, zones_bytes = decode_upload(files["zones"], {".geojson"})
        rules_name, rules_bytes = decode_upload(files["rules"], {".json"})
        if len({observations_name.casefold(), zones_name.casefold(), rules_name.casefold()}) != 3:
            raise ApiError("Risk-rule uploaded file names must be unique.")
        run_id = make_run_id("risk")
        raw_root = self.root / "shared" / "data" / "raw" / "07-risk-rule-engine" / "runs" / run_id
        processed_root = self.root / "shared" / "data" / "processed" / "07-risk-rule-engine" / run_id
        raw_root.mkdir(parents=True, exist_ok=False)
        processed_root.mkdir(parents=True, exist_ok=False)
        staged = ((observations_name, observations_bytes), (zones_name, zones_bytes), (rules_name, rules_bytes))
        for name, content in staged:
            (raw_root / name).write_bytes(content)
            (processed_root / name).write_bytes(content)
        output = self.root / "shared" / "outputs" / "07-risk-rule-engine" / "runs" / run_id
        python = self.root / ".venvs" / "07-risk-rule-engine" / "Scripts" / "python.exe"
        if not python.is_file():
            raise ApiError("Risk-rule virtual environment is unavailable. Run the capability setup first.")
        command = [
            str(python), str(self.root / "capabilities" / "07-risk-rule-engine" / "run_risk_rule_engine.py"),
            "--observations", str(processed_root / observations_name), "--zones", str(processed_root / zones_name),
            "--rules", str(processed_root / rules_name), "--output", str(output),
        ]
        with RUN_LOCK:
            self.run_command(command, 600)
        metadata_path = output / "run_metadata.json"
        if not metadata_path.is_file():
            raise ApiError("Risk-rule script finished without the expected result metadata.")
        metadata = load_json(metadata_path)
        metadata["input_dir"] = relative_path(self.root, processed_root)
        metadata["raw_input_dir"] = relative_path(self.root, raw_root)
        metadata["source_bytes"] = {name: len(content) for name, content in staged}
        metadata_path.write_text(json.dumps(metadata, ensure_ascii=False, indent=2), encoding="utf-8")
        return next(item for item in risk_rule_runs(self.root) if item["id"] == run_id)
    def send_json(self, status: HTTPStatus, payload: dict[str, Any]) -> None:
        body = json.dumps(payload, ensure_ascii=False).encode("utf-8")
        self.send_response(status)