Refactor ENVI D-InSAR production runs

This commit is contained in:
2026-04-16 17:09:24 +08:00
parent 5bdcaf74df
commit d83cdfea52
12 changed files with 1796 additions and 33 deletions
+603 -10
View File
@@ -5,10 +5,12 @@ import io
import json
import logging
import os
import re
import subprocess
import sys
import tempfile
import time
import uuid
from typing import Callable, Awaitable, Optional, Dict, Any, List
logger = logging.getLogger(__name__)
@@ -21,8 +23,11 @@ from ..models import SystemJobORM, DinsarResultORM, HazardPointORM, DinsarTaskIt
from ..scheduler import scan_data_job
from .data_service import data_service
from .dinsar_compat_service import dinsar_compat_service
from .dinsar_naming import build_run_key
from .dinsar_production_service import dinsar_production_service
from .dinsar_read_service import dinsar_read_service
from .dinsar_scan_service import dinsar_scan_service
from .engine_lock_service import engine_lock_service
from .psinsar_catalog_service import psinsar_catalog_service
from .result_catalog_service import result_catalog_service
from .task_service import task_service
@@ -924,6 +929,258 @@ def _kill_process_tree(pid: int) -> None:
print(f"[WARN] _kill_process_tree: pid={pid}{exc}")
_ENVI_RESULT_NAME_RE = re.compile(r"^.+_rsp_disp$", re.IGNORECASE)
def _format_envi_keepalive(progress: Optional[Dict[str, Any]]) -> tuple[str, Optional[int]]:
if not progress:
return "ENVI processing...", None
step = progress.get("step", 0)
total = progress.get("total_steps", 6)
step_msg = progress.get("message", "")
pair_index = progress.get("pair_index", 0)
total_pairs = progress.get("total_pairs", 0)
pair_name = progress.get("pair_name", "")
step_part = f"Step {step}/{total}: {step_msg}" if step_msg else f"Step {step}/{total}"
if total_pairs > 0 and pair_index > 0:
pair_part = f"Pair {pair_index}/{total_pairs}"
if pair_name:
pair_part += f" ({pair_name})"
message = f"{pair_part} | {step_part}"
else:
message = step_part
progress_value: Optional[int] = None
if isinstance(step, (int, float)) and isinstance(total, (int, float)) and total > 0:
if total_pairs > 0 and pair_index > 0:
pair_frac = (pair_index - 1 + step / total) / total_pairs
progress_value = min(90, 10 + int(80 * pair_frac))
else:
progress_value = min(90, 10 + int(80 * step / total))
return message, progress_value
def _clear_envi_progress_file(job_id: Optional[str]) -> None:
if not job_id:
return
try:
progress_path = _get_envi_progress_file(job_id)
if os.path.isfile(progress_path):
os.remove(progress_path)
except OSError:
pass
def _find_latest_envi_result(output_dir: str) -> Dict[str, Any]:
matches: List[tuple[float, str]] = []
try:
for entry in os.scandir(output_dir):
if not entry.is_file():
continue
if entry.name.lower().endswith((".hdr", ".sml")):
continue
if not _ENVI_RESULT_NAME_RE.match(entry.name):
continue
try:
stat = entry.stat()
matches.append((max(stat.st_mtime, stat.st_ctime), entry.path))
except OSError:
matches.append((0.0, entry.path))
except OSError as exc:
raise RuntimeError(f"Failed to scan ENVI output directory: {output_dir}: {exc}") from exc
if not matches:
raise RuntimeError(f"No ENVI displacement result found under: {output_dir}")
matches.sort(key=lambda item: (item[0], item[1]), reverse=True)
primary_file = matches[0][1]
source_files = [primary_file]
for ext in (".hdr", ".sml"):
sidecar = primary_file + ext
if os.path.isfile(sidecar):
source_files.append(sidecar)
return {
"primary_file": primary_file,
"source_files": source_files,
}
async def _run_envi_runner_command(
job: SystemJobORM,
runner_cmd: List[str],
*,
absolute_timeout_seconds: int,
keepalive_formatter: Optional[Callable[[Optional[Dict[str, Any]]], tuple[str, Optional[int]]]] = None,
register_pid: Optional[Callable[[int], Awaitable[None]]] = None,
) -> Dict[str, Any]:
if not job.task_id:
raise ValueError(f"{job.job_type} requires task_id for progress tracking.")
formatter = keepalive_formatter or _format_envi_keepalive
progress_file = _get_envi_progress_file(job.job_id)
pid_ready = asyncio.Event()
proc_state: Dict[str, Any] = {}
loop = asyncio.get_running_loop()
async def _task_keepalive():
while True:
await asyncio.sleep(30)
try:
progress = _read_progress_file(progress_file)
message, progress_value = formatter(progress)
await task_service.update_task(job.task_id, message=message, progress=progress_value)
except Exception as keepalive_exc:
print(f"[keepalive] WARNING: failed to update task {job.task_id}: {keepalive_exc}")
def _run_with_monitoring() -> Dict[str, Any]:
stdout_fd = None
stderr_fd = None
stdout_path = None
stderr_path = None
proc = None
try:
stdout_fd, stdout_path = tempfile.mkstemp(suffix="_stdout.txt")
stderr_fd, stderr_path = tempfile.mkstemp(suffix="_stderr.txt")
proc = subprocess.Popen(
runner_cmd,
stdout=stdout_fd,
stderr=stderr_fd,
cwd=type(settings).PROJECT_ROOT,
)
proc_state["pid"] = proc.pid
loop.call_soon_threadsafe(pid_ready.set)
os.close(stdout_fd)
os.close(stderr_fd)
absolute_start = time.time()
last_activity = time.time()
last_step_msg = ""
while proc.poll() is None:
time.sleep(_ENVI_MONITOR_INTERVAL)
now = time.time()
progress = _read_progress_file(progress_file)
if progress:
ts = progress.get("timestamp", 0)
if ts > last_activity:
last_activity = ts
msg = progress.get("message", "")
if msg and msg != last_step_msg:
last_step_msg = msg
output_dir = ""
if progress and progress.get("output_dir"):
output_dir = progress["output_dir"]
if output_dir:
dir_mtime = _scan_latest_mtime(output_dir)
if dir_mtime and dir_mtime > last_activity:
last_activity = dir_mtime
if (now - last_activity) > _ENVI_FILE_STALE_SECONDS:
_kill_process_tree(proc.pid)
try:
proc.wait(timeout=15)
except Exception as exc:
print(f"[WARN] stale kill: proc.wait timeout — {exc}")
raise RuntimeError(
f"ENVI subprocess stale: no file activity for "
f"{int(now - last_activity)}s "
f"(threshold={_ENVI_FILE_STALE_SECONDS}s). "
f"Last step: {last_step_msg}"
)
if (now - absolute_start) > absolute_timeout_seconds:
_kill_process_tree(proc.pid)
try:
proc.wait(timeout=15)
except Exception:
pass
raise RuntimeError(
f"ENVI subprocess exceeded absolute timeout: "
f"{int(now - absolute_start)}s > {absolute_timeout_seconds}s. "
f"Last step: {last_step_msg}"
)
output_dir = ""
progress = _read_progress_file(progress_file)
if progress and progress.get("output_dir"):
output_dir = progress["output_dir"]
if output_dir:
stable_count = 0
prev_sizes: Dict[str, int] = {}
wait_start = time.time()
while stable_count < _ENVI_STABILITY_ROUNDS:
if (time.time() - wait_start) > _ENVI_STABILITY_MAX_WAIT:
break
time.sleep(_ENVI_STABILITY_CHECK_INTERVAL)
cur_sizes = _collect_file_sizes(output_dir)
if cur_sizes == prev_sizes:
stable_count += 1
else:
stable_count = 0
prev_sizes = cur_sizes
with open(stdout_path, "r", encoding="utf-8", errors="replace") as fp:
stdout = fp.read()
with open(stderr_path, "r", encoding="utf-8", errors="replace") as fp:
stderr = fp.read()
finally:
for fd in (stdout_fd, stderr_fd):
if isinstance(fd, int):
try:
os.close(fd)
except OSError:
pass
for path in (stdout_path, stderr_path):
if path:
try:
os.unlink(path)
except OSError:
pass
if proc is None:
raise RuntimeError("ENVI runner subprocess failed to start.")
if proc.returncode != 0:
raise RuntimeError(
"ENVI runner subprocess failed. "
f"rc={proc.returncode} "
f"stderr={(stderr or '').strip()[:1200]!r}"
)
output = (stdout or "").strip()
if not output:
raise RuntimeError("ENVI runner subprocess returned empty output.")
try:
return json.loads(output.splitlines()[-1])
except Exception as exc:
raise RuntimeError(f"ENVI runner returned non-JSON output: {output[:1200]!r}") from exc
keepalive_task = asyncio.create_task(_task_keepalive())
runner_task = asyncio.create_task(asyncio.to_thread(_run_with_monitoring))
try:
if register_pid is not None:
try:
await asyncio.wait_for(pid_ready.wait(), timeout=30)
pid_value = int(proc_state.get("pid") or 0)
if pid_value > 0:
await register_pid(pid_value)
except asyncio.TimeoutError:
pass
return await runner_task
finally:
keepalive_task.cancel()
try:
await keepalive_task
except asyncio.CancelledError:
pass
async def _run_envi_workflow_job(
job: SystemJobORM,
workflow: str,
@@ -1215,23 +1472,359 @@ async def _run_envi_workflow_job(
)
async def _run_dinsar_production_controller(job: SystemJobORM) -> None:
if not job.task_id:
raise ValueError("IDL_RUN_DINSAR production controller requires task_id.")
payload = job.payload or {}
production_run_id = str(payload.get("production_run_id") or "").strip()
if not production_run_id:
raise ValueError("IDL_RUN_DINSAR production controller requires production_run_id.")
async with AsyncSessionLocal() as db:
run = await dinsar_production_service.get_run(production_run_id, db)
if run is None:
raise ValueError(f"D-InSAR production run not found: {production_run_id}")
items = await dinsar_production_service.list_run_items(run.run_id, db)
if not items:
raise ValueError(f"D-InSAR production run has no items: {production_run_id}")
workflow = "dinsar_custom" if str(run.mode or "").strip().lower() == "custom" else "dinsar"
params = run.params_json or {}
timeout_seconds_raw = params.get("timeout_seconds")
timeout_seconds = int(timeout_seconds_raw) if timeout_seconds_raw not in (None, "") else None
absolute_timeout_seconds = max(_ENVI_PER_TASK_TIMEOUT, int(timeout_seconds or 0)) if timeout_seconds else _ENVI_PER_TASK_TIMEOUT
total_items = len(items)
run_log = dinsar_production_service.append_run_log
async def _refresh_cancel_state() -> bool:
await db.refresh(run)
current_task = await task_service.get_task(job.task_id)
task_cancelled = bool(current_task and current_task.status == "CANCELLED")
if task_cancelled and not run.cancel_requested:
run.cancel_requested = True
await db.commit()
return bool(run.cancel_requested or task_cancelled)
await task_service.start_task(
job.task_id,
message=f"Starting ENVI D-InSAR production run {run.run_id} ({total_items} items)...",
)
await task_service.add_log(
job.task_id,
"INFO",
(
f"D-InSAR production controller started. run_id={run.run_id} "
f"profile={run.profile_code} mode={run.mode} items={total_items}"
),
)
await dinsar_production_service.mark_run_started(
run,
db=db,
message=f"Preparing {total_items} D-InSAR item(s)",
)
run_log(run.run_id, f"[start] workflow={workflow} items={total_items} source_root={run.source_root}")
successful_output_dirs: List[str] = []
async with engine_lock_service.acquire("envi_taskengine"):
await task_service.update_task(
job.task_id,
progress=5,
message=f"ENVI engine acquired. Preparing {total_items} item(s)...",
)
for item_index, item in enumerate(items, start=1):
if await _refresh_cancel_state():
await task_service.add_log(
job.task_id,
"WARNING",
f"Cancellation detected before item {item.task_alias or item.task_name}.",
)
break
await db.refresh(item)
if str(item.status or "").upper() in {"COMPLETED", "FAILED", "SKIPPED", "CANCELLED"}:
if item.status == "COMPLETED" and item.latest_output_dir and os.path.isdir(item.latest_output_dir):
successful_output_dirs.append(item.latest_output_dir)
continue
run_key = (
f"{build_run_key('sarscape', run.profile_code, started_at=datetime.utcnow())}"
f"_{item.id}_{uuid.uuid4().hex[:6]}"
)
execution = await dinsar_production_service.begin_item_execution(
run=run,
item=item,
run_key=run_key,
db=db,
)
runner_cmd = [
sys.executable,
"-m",
"backend.app.services.envi_runner_cli",
"--workflow",
workflow,
"--task-dir",
str(item.source_task_dir),
"--output-dir",
str(execution.output_dir),
"--source-root",
str(run.source_root),
"--job-id",
str(job.job_id),
"--run-key",
str(run_key),
"--profile-code",
str(run.profile_code),
]
if timeout_seconds is not None:
runner_cmd.extend(["--timeout-seconds", str(timeout_seconds)])
def _keepalive_formatter(progress: Optional[Dict[str, Any]]) -> tuple[str, Optional[int]]:
message, progress_value = _format_envi_keepalive(progress)
prefix = f"{item_index}/{total_items} {item.task_alias or item.task_name}: "
if progress_value is None:
return prefix + message, None
local_fraction = max(0.0, min(1.0, progress_value / 100.0))
overall = ((item_index - 1) + local_fraction) / max(1, total_items)
return prefix + message, min(98, max(1, int(overall * 100)))
async def _register_pid(pid: int) -> None:
async with AsyncSessionLocal() as pid_db:
await dinsar_production_service.set_execution_pid(
execution.execution_id,
pid,
db=pid_db,
)
await task_service.add_log(
job.task_id,
"INFO",
(
f"[{item_index}/{total_items}] Launching {item.task_alias or item.task_name} "
f"-> {execution.output_dir}"
),
)
run_log(
run.run_id,
(
f"[item-start] {item_index}/{total_items} "
f"{item.task_alias or item.task_name} run_key={run_key} output={execution.output_dir}"
),
)
try:
run_meta = await _run_envi_runner_command(
job,
runner_cmd,
absolute_timeout_seconds=absolute_timeout_seconds,
keepalive_formatter=_keepalive_formatter,
register_pid=_register_pid,
)
result_files = await asyncio.to_thread(_find_latest_envi_result, execution.output_dir)
metrics = {
"duration_seconds": run_meta.get("duration_seconds"),
"summary": run_meta.get("summary") or {},
}
if run_meta.get("task_results"):
metrics["task_result"] = run_meta["task_results"][0]
manifest_path = await asyncio.to_thread(
dinsar_production_service.build_execution_manifest,
run=run,
item=item,
execution=execution,
primary_file=result_files["primary_file"],
source_files=result_files["source_files"],
metrics=metrics,
)
await asyncio.to_thread(
dinsar_production_service.write_current_pointer,
item=item,
execution=execution,
manifest_path=manifest_path,
primary_file=result_files["primary_file"],
source_files=result_files["source_files"],
)
await dinsar_production_service.mark_item_completed(
run=run,
item=item,
execution=execution,
manifest_path=manifest_path,
metrics=metrics,
db=db,
)
successful_output_dirs.append(execution.output_dir)
await task_service.add_log(
job.task_id,
"INFO",
f"[{item_index}/{total_items}] Completed {item.task_alias or item.task_name}",
)
run_log(
run.run_id,
f"[item-ok] {item_index}/{total_items} {item.task_alias or item.task_name}",
)
except Exception as exc:
cancelled = await _refresh_cancel_state()
if cancelled:
await dinsar_production_service.mark_item_cancelled(
run=run,
item=item,
execution=execution,
error_message=f"Cancelled while processing {item.task_alias or item.task_name}",
db=db,
)
await task_service.add_log(
job.task_id,
"WARNING",
f"[{item_index}/{total_items}] Cancelled {item.task_alias or item.task_name}: {exc}",
)
run_log(
run.run_id,
f"[item-cancelled] {item_index}/{total_items} {item.task_alias or item.task_name}: {exc}",
)
break
error_message = str(exc)
await dinsar_production_service.mark_item_failed(
run=run,
item=item,
execution=execution,
error_message=error_message,
db=db,
)
await task_service.add_log(
job.task_id,
"WARNING",
f"[{item_index}/{total_items}] Failed {item.task_alias or item.task_name}: {error_message}",
)
run_log(
run.run_id,
f"[item-failed] {item_index}/{total_items} {item.task_alias or item.task_name}: {error_message}",
)
finally:
_clear_envi_progress_file(job.job_id)
publish_result = None
publish_error = None
publish_dirs = _dedupe_existing_dirs(successful_output_dirs)
if publish_dirs:
await task_service.update_task(
job.task_id,
progress=99,
message=f"Publishing {len(publish_dirs)} successful result package(s)...",
)
try:
publish_result = await result_catalog_service.publish_from_sources(db, publish_dirs)
await task_service.add_log(
job.task_id,
"INFO",
f"Published {publish_result.get('processed', 0)} result package(s).",
)
run_log(
run.run_id,
f"[publish] processed={publish_result.get('processed', 0)} failed={publish_result.get('failed', 0)}",
)
except Exception as exc:
publish_error = str(exc)
await task_service.add_log(
job.task_id,
"WARNING",
f"Result catalog publish failed: {publish_error}",
)
run_log(run.run_id, f"[publish-failed] {publish_error}")
cancelled = await _refresh_cancel_state()
await db.refresh(run)
latest_message = ""
final_status = "COMPLETED"
if publish_error:
final_status = "FAILED"
latest_message = f"Result catalog publish failed: {publish_error}"
elif cancelled:
final_status = "CANCELLED"
latest_message = (
f"D-InSAR production cancelled. completed={run.completed_items} "
f"failed={run.failed_items} total={run.total_items}"
)
elif int(run.failed_items or 0) > 0:
final_status = "FAILED"
latest_message = (
f"D-InSAR production finished with failures. completed={run.completed_items} "
f"failed={run.failed_items} total={run.total_items}"
)
else:
latest_message = (
f"D-InSAR production completed. completed={run.completed_items} "
f"failed={run.failed_items} total={run.total_items}"
)
summary_payload = {
"workflow": workflow,
"engine_code": run.engine_code,
"profile_code": run.profile_code,
"mode": run.mode,
"total_items": run.total_items,
"completed_items": run.completed_items,
"failed_items": run.failed_items,
"skipped_items": run.skipped_items,
"publish": publish_result,
"publish_error": publish_error,
"published_output_dirs": publish_dirs,
}
await dinsar_production_service.finalize_run(
run,
db=db,
status=final_status,
summary_payload=summary_payload,
latest_message=latest_message,
)
run_log(run.run_id, f"[finish] status={final_status} message={latest_message}")
if final_status == "COMPLETED":
await task_service.update_task(
job.task_id,
status="COMPLETED",
progress=100,
message=latest_message,
)
return
task_status = "CANCELLED" if final_status == "CANCELLED" else "FAILED"
await task_service.update_task(
job.task_id,
status=task_status,
progress=100,
message=latest_message,
)
raise RuntimeError(latest_message)
async def _handle_idl_run_import(job: SystemJobORM) -> None:
await _run_envi_workflow_job(
job,
workflow="import",
start_message="Starting ENVI Import workflow...",
)
async with engine_lock_service.acquire("envi_taskengine"):
await _run_envi_workflow_job(
job,
workflow="import",
start_message="Starting ENVI Import workflow...",
)
async def _handle_idl_run_dinsar(job: SystemJobORM) -> None:
payload = job.payload or {}
production_run_id = str(payload.get("production_run_id") or "").strip()
if production_run_id:
await _run_dinsar_production_controller(job)
return
mode = payload.get("mode", "metatask")
workflow = "dinsar_custom" if mode == "custom" else "dinsar"
await _run_envi_workflow_job(
job,
workflow=workflow,
start_message=f"Starting ENVI D-InSAR workflow (mode={mode})...",
)
async with engine_lock_service.acquire("envi_taskengine"):
await _run_envi_workflow_job(
job,
workflow=workflow,
start_message=f"Starting ENVI D-InSAR workflow (mode={mode})...",
)
async def _handle_isce2_run(job: SystemJobORM) -> None: