From 01776511b66bfd87b4f8ef57d3fefc1d99e1e8f6 Mon Sep 17 00:00:00 2001
From: shuishen <1109946754@qq.com>
Date: Wed, 19 Aug 2026 16:46:52 +0800
Subject: [PATCH] feat:变化检测相关优化调整
---
scripts/serve_workbench_console.py | 335 +++++++++++++++++++++++++++++++++++++++++++++++++++++++
1 files changed, 333 insertions(+), 2 deletions(-)
diff --git a/scripts/serve_workbench_console.py b/scripts/serve_workbench_console.py
index e6fd1b1..9219509 100644
--- a/scripts/serve_workbench_console.py
+++ b/scripts/serve_workbench_console.py
@@ -38,6 +38,12 @@
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",
@@ -47,7 +53,11 @@
)
SAFE_FILE_NAME = re.compile(r"[^A-Za-z0-9._-]+")
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()
class ApiError(ValueError):
@@ -153,7 +163,7 @@
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 not isinstance(artifacts, dict):
+ 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
@@ -187,6 +197,69 @@
}
)
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]]:
@@ -255,6 +328,18 @@
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
@@ -283,6 +368,13 @@
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)})
@@ -354,7 +446,11 @@
role = query.get("role", [""])[0]
if role not in {"before", "after"}:
raise ApiError("Change-detection upload role must be before or after.")
- name = self.headers.get("X-Upload-Name", "")
+ encoded_name = self.headers.get("X-Upload-Name", "")
+ try:
+ name = unquote(encoded_name)
+ except Exception as exc:
+ raise ApiError("The uploaded filename is invalid.") from exc
safe_name = safe_file_name(name, {".jpg", ".jpeg", ".png", ".tif", ".tiff"})
content_length = self.headers.get("Content-Length")
if content_length is None or not content_length.isdigit():
@@ -528,6 +624,241 @@
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"),
+ ("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 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),
+ "artifacts": {
+ "probability_raster": "change_probability.tif",
+ "raw_mask_raster": "change_mask.tif",
+ "mask_raster": "change_mask.tif",
+ "overlay": "change_overlay.jpg",
+ "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)
--
Gitblit v1.9.3