From fced6f4a7f7f9b553c4236eaeed27da86f5cf86b Mon Sep 17 00:00:00 2001 From: Harmon Date: Fri, 26 Jun 2026 11:27:05 +0800 Subject: [PATCH] Add LandSAR cluster data transport --- .env.example | 7 + backend/app/auth_service.py | 2 +- backend/app/main.py | 2 + backend/app/routers/cluster.py | 282 +++++++++++++++++ backend/app/services/cluster_transport.py | 290 ++++++++++++++++++ backend/app/services/health_service.py | 48 +++ backend/app/services/job_handlers.py | 71 +++-- config/landsar_cluster_worker.env.example | 21 ++ docs/INDEX.md | 2 + ..._CLUSTER_DATA_TRANSPORT_DESIGN_20260625.md | 235 ++++++++++++++ scripts/start_landsar_cluster_worker.ps1 | 161 ++++++++++ .../test_landsar_cluster_worker_remote.ps1 | 157 ++++++++++ 12 files changed, 1258 insertions(+), 20 deletions(-) create mode 100644 backend/app/routers/cluster.py create mode 100644 backend/app/services/cluster_transport.py create mode 100644 docs/LANDSAR_CLUSTER_DATA_TRANSPORT_DESIGN_20260625.md create mode 100644 scripts/test_landsar_cluster_worker_remote.ps1 diff --git a/.env.example b/.env.example index 99a3c58..12dd076 100644 --- a/.env.example +++ b/.env.example @@ -342,6 +342,13 @@ JOB_WORKER_HEARTBEAT_INTERVAL=5 # Main server only. Comma/semicolon-separated IPv4 addresses or CIDR blocks allowed # to run LandSAR cluster workers against this PostgreSQL server. LANDSAR_CLUSTER_ALLOWED_WORKER_IPS=192.168.1.6 +# Main server and remote workers must share the same non-empty token for +# /api/cluster input download and result upload. +CLUSTER_SHARED_TOKEN= +CLUSTER_TRANSFER_TIMEOUT_SECONDS=3600 +CLUSTER_MATERIALIZE_TEMP_DIR=D:\Task_Pool\_cluster_temp +# Remote LandSAR workers should point this to the main backend URL. +CLUSTER_MAIN_SERVER_URL= # Empty means the worker can claim all job types. Remote LandSAR nodes should set: # JOB_WORKER_ALLOWED_TYPES=LANDSAR_CLUSTER_ITEM JOB_WORKER_ALLOWED_TYPES= diff --git a/backend/app/auth_service.py b/backend/app/auth_service.py index a7f9c32..984fa85 100644 --- a/backend/app/auth_service.py +++ b/backend/app/auth_service.py @@ -25,7 +25,7 @@ VALID_ROLES = {ROLE_ADMIN, ROLE_VIEWER} SESSION_COOKIE_NAME = settings.AUTH_SESSION_COOKIE_NAME SESSION_TTL_HOURS = read_int_env("SESSION_TTL_HOURS", 12) -COOKIE_SECURE = read_bool_env("AUTH_COOKIE_SECURE", True) +COOKIE_SECURE = read_bool_env("AUTH_COOKIE_SECURE", False) COOKIE_SAMESITE = (settings.AUTH_COOKIE_SAMESITE or "lax").strip().lower() if COOKIE_SAMESITE not in {"lax", "strict", "none"}: COOKIE_SAMESITE = "lax" diff --git a/backend/app/main.py b/backend/app/main.py index 1c38e6e..95980a3 100644 --- a/backend/app/main.py +++ b/backend/app/main.py @@ -11,6 +11,7 @@ ensure_project_env_loaded() from . import database from .api import router as api_router +from .routers.cluster import router as cluster_router from .db_maintenance import ensure_database_ready from .scheduler import scheduler_manager from .services.health_service import get_health_status @@ -301,6 +302,7 @@ app.add_middleware( allow_headers=["*"], ) +app.include_router(cluster_router, prefix="/api") app.include_router(api_router, prefix="/api") diff --git a/backend/app/routers/cluster.py b/backend/app/routers/cluster.py new file mode 100644 index 0000000..896f83d --- /dev/null +++ b/backend/app/routers/cluster.py @@ -0,0 +1,282 @@ +""" +Cluster data-transport endpoints for LandSAR distributed processing. + +Provides HTTP-based input materialize (download) and result upload +so that remote Windows workers can pull Task_Pool pair data and +push finished products back to the main server without SMB / UNC. +""" +from __future__ import annotations + +import os +import shutil +import tempfile +import zipfile +from typing import Optional + +from fastapi import ( + APIRouter, + Depends, + File, + Form, + Header, + HTTPException, + UploadFile, +) +from fastapi.responses import FileResponse +from starlette.background import BackgroundTask +from sqlalchemy.ext.asyncio import AsyncSession + +from ..config import settings +from ..database import get_db +from ..models.orm import DinsarProductionRunItemORM +from ..services.cluster_transport import safe_extract_zip + +router = APIRouter() + + +def _cluster_materialize_temp_dir() -> str: + explicit = str(os.environ.get("CLUSTER_MATERIALIZE_TEMP_DIR") or "").strip() + if not explicit: + explicit = str( + getattr(settings, "CLUSTER_MATERIALIZE_TEMP_DIR", "") or "" + ).strip() + if explicit: + return os.path.normpath(explicit) + task_pool = str(getattr(settings, "TASK_POOL_ROOT", "") or "").strip() + if task_pool: + return os.path.normpath(os.path.join(task_pool, "_cluster_temp")) + return os.path.normpath( + os.path.join( + os.path.dirname(__file__), "..", "..", "runtime", "_cluster_temp" + ) + ) + + +def _cluster_shared_token() -> str: + return str(os.environ.get("CLUSTER_SHARED_TOKEN") or "").strip() + + +def _require_cluster_token( + x_cluster_token: Optional[str] = Header(default=None), +) -> None: + expected = _cluster_shared_token() + if not expected: + raise HTTPException( + status_code=503, + detail="CLUSTER_SHARED_TOKEN is not configured.", + ) + if x_cluster_token != expected: + raise HTTPException(status_code=403, detail="Invalid cluster token.") + + +# --------------------------------------------------------------------------- +# Download input package +# --------------------------------------------------------------------------- + +@router.get("/cluster/input-package/{item_id}") +async def download_cluster_input_package( + item_id: int, + _cluster_token: None = Depends(_require_cluster_token), + db: AsyncSession = Depends(get_db), +): + """Package a cluster item's source-task directory as a zip and stream it. + + The remote worker calls this when *source_task_dir* does not exist + locally. The zip mirrors the Task_Pool layout: + + Task_YYYYMMDD_YYYYMMDD/ + master/ ... + slave/ ... + orbit/ ... + pair_metadata.json + """ + item = await db.get(DinsarProductionRunItemORM, int(item_id)) + if item is None: + raise HTTPException(status_code=404, detail="Cluster item not found.") + + source_dir = os.path.normpath(str(item.source_task_dir or "")) + if not source_dir or not os.path.isdir(source_dir): + raise HTTPException( + status_code=404, + detail=f"Source task directory not found: {source_dir}", + ) + + task_name = os.path.basename(source_dir) + parent_dir = os.path.dirname(source_dir) + temp_root = _cluster_materialize_temp_dir() + os.makedirs(temp_root, exist_ok=True) + + package_root = tempfile.mkdtemp( + prefix=f"cluster_input_{item_id}_", + dir=temp_root, + ) + tmp_base = os.path.join(package_root, task_name) + try: + zip_path = shutil.make_archive( + tmp_base, + "zip", + root_dir=parent_dir, + base_dir=task_name, + ) + except Exception as exc: + try: + shutil.rmtree(package_root) + except Exception: + pass + raise HTTPException( + status_code=500, + detail=f"Failed to create input package: {exc}", + ) + + def _cleanup(): + try: + if os.path.isdir(package_root): + shutil.rmtree(package_root) + except Exception: + pass + + return FileResponse( + path=zip_path, + media_type="application/zip", + filename=f"{task_name}.zip", + background=BackgroundTask(_cleanup), + ) + + +# --------------------------------------------------------------------------- +# Upload result package +# --------------------------------------------------------------------------- + +@router.post("/cluster/upload-result/{item_id}") +async def upload_cluster_result( + item_id: int, + run_id: str = Form(...), + run_key: str = Form(...), + result_zip: UploadFile = File(...), + _cluster_token: None = Depends(_require_cluster_token), + db: AsyncSession = Depends(get_db), +): + """Receive a result zip from a cluster worker and register it in the + D-InSAR catalog. + + The zip is extracted under the cluster item's standard + ``results_root_dir/runs/`` directory and then + *result_catalog_service.publish_from_sources* is called so the product + is immediately visible on the main server. + """ + item = await db.get(DinsarProductionRunItemORM, int(item_id)) + if item is None: + raise HTTPException(status_code=404, detail="Cluster item not found.") + + normalized_run_id = str(run_id or "").strip() + if normalized_run_id != str(item.run_id or "").strip(): + raise HTTPException( + status_code=400, + detail="run_id does not match item.", + ) + + publish_root = os.path.normpath( + str(getattr(settings, "DINSAR_PRODUCT_DIR", "") or "") + ) + if not publish_root or not os.path.isdir(publish_root): + raise HTTPException( + status_code=500, + detail="DINSAR_PRODUCT_DIR is not configured.", + ) + + normalized_run_key = str(run_key or "").strip() + if not normalized_run_key: + raise HTTPException(status_code=400, detail="run_key is required.") + + item_results_root = os.path.normpath(str(item.results_root_dir or "")) + if not item_results_root: + raise HTTPException( + status_code=500, + detail="Cluster item results_root_dir is not configured.", + ) + publish_root_abs = os.path.abspath(publish_root) + item_results_root_abs = os.path.abspath(item_results_root) + if item_results_root_abs != publish_root_abs and not item_results_root_abs.startswith( + publish_root_abs + os.sep + ): + raise HTTPException( + status_code=400, + detail="Cluster item results_root_dir is outside DINSAR_PRODUCT_DIR.", + ) + + extract_dir = os.path.join(item_results_root_abs, "runs", normalized_run_key) + os.makedirs(os.path.dirname(extract_dir), exist_ok=True) + + extract_parent = os.path.dirname(extract_dir) + os.makedirs(extract_parent, exist_ok=True) + upload_root = tempfile.mkdtemp( + prefix=f"cluster_upload_{item_id}_", + dir=extract_parent, + ) + tmp_extract = os.path.join(upload_root, "extract") + tmp_zip = os.path.join(upload_root, "upload.zip") + backup_dir: Optional[str] = None + try: + with open(tmp_zip, "wb") as fh: + while chunk := await result_zip.read(8 * 1024 * 1024): # 8 MiB + fh.write(chunk) + + os.makedirs(tmp_extract, exist_ok=True) + with zipfile.ZipFile(tmp_zip, "r") as zf: + safe_extract_zip(zf, tmp_extract) + + if os.path.isdir(extract_dir): + backup_dir = f"{extract_dir}._replace_backup" + if os.path.isdir(backup_dir): + shutil.rmtree(backup_dir) + os.replace(extract_dir, backup_dir) + os.replace(tmp_extract, extract_dir) + if backup_dir and os.path.isdir(backup_dir): + shutil.rmtree(backup_dir) + backup_dir = None + except zipfile.BadZipFile: + raise HTTPException( + status_code=400, + detail="Uploaded file is not a valid zip archive.", + ) + except ValueError as exc: + raise HTTPException(status_code=400, detail=str(exc)) from exc + except Exception: + if backup_dir and os.path.isdir(backup_dir) and not os.path.exists(extract_dir): + try: + os.replace(backup_dir, extract_dir) + backup_dir = None + except Exception: + pass + raise + finally: + try: + if os.path.isdir(upload_root): + shutil.rmtree(upload_root) + except Exception: + pass + try: + if backup_dir and os.path.isdir(backup_dir): + shutil.rmtree(backup_dir) + except Exception: + pass + + # ----- catalog registration ------------------------------------------------ + from ..services.result_catalog_service import result_catalog_service as rcs + + try: + publish_result = await rcs.publish_from_sources(db, [extract_dir]) + processed = int(publish_result.get("processed", 0) or 0) + if processed > 0: + await rcs.rebuild_catalog(db, full_rebuild=True) + return { + "registered": processed > 0, + "processed": processed, + "failed": int(publish_result.get("failed", 0) or 0), + "catalog_path": extract_dir, + } + except Exception as exc: + raise HTTPException( + status_code=500, + detail=f"Catalog registration failed: {exc}", + ) from exc diff --git a/backend/app/services/cluster_transport.py b/backend/app/services/cluster_transport.py new file mode 100644 index 0000000..14ed155 --- /dev/null +++ b/backend/app/services/cluster_transport.py @@ -0,0 +1,290 @@ +""" +Cluster data-transport helpers for LandSAR distributed processing. + +Used by _handle_landsar_cluster_item in job_handlers.py to pull input +data from the main server and push results back via HTTP. +""" +from __future__ import annotations + +import io +import json +import os +import shutil +import tempfile +import urllib.error +import urllib.parse +import urllib.request +import zipfile +from typing import TYPE_CHECKING + +if TYPE_CHECKING: + from ..models.orm import DinsarProductionRunItemORM, DinsarProductionRunORM + + +def _read_cluster_env(name: str) -> str: + return str(os.environ.get(name) or "").strip() + + +def _cluster_transfer_timeout() -> int: + try: + from ..config import read_int_env + + return read_int_env( + "CLUSTER_TRANSFER_TIMEOUT_SECONDS", + 3600, + minimum=60, + maximum=86400, + ) + except Exception: + return 3600 + + +def _cluster_request_headers() -> dict[str, str]: + token = _read_cluster_env("CLUSTER_SHARED_TOKEN") + if not token: + raise RuntimeError("CLUSTER_SHARED_TOKEN is not configured.") + return {"X-Cluster-Token": token} + + +def safe_extract_zip(zf: zipfile.ZipFile, target_dir: str) -> None: + """Extract a zip after verifying all members stay inside target_dir.""" + target_root = os.path.abspath(target_dir) + for member in zf.infolist(): + member_name = member.filename.replace("\\", "/") + if ( + not member_name + or member_name.startswith("/") + or os.path.splitdrive(member_name)[0] + or member_name.startswith("../") + or "/../" in f"/{member_name}/" + ): + raise ValueError(f"Unsafe zip member path: {member.filename}") + destination = os.path.abspath(os.path.join(target_root, member_name)) + if destination != target_root and not destination.startswith( + target_root + os.sep + ): + raise ValueError(f"Unsafe zip member path: {member.filename}") + zf.extractall(target_root) + + +def _resolve_cluster_server_url() -> str: + """Return the main-server HTTP base URL for cluster data transport. + + Prefers CLUSTER_MAIN_SERVER_URL; falls back to the DATABASE_URL host. + Returns ``http://127.0.0.1`` when nothing is configured (main-server / + local worker). + """ + from ..config import settings + + explicit = _read_cluster_env("CLUSTER_MAIN_SERVER_URL") or str( + getattr(settings, "CLUSTER_MAIN_SERVER_URL", "") or "" + ).strip() + if explicit: + return explicit.rstrip("/") + + db_url = str(getattr(settings, "DATABASE_URL", "") or "") + if "@" in db_url: + host_part = db_url.split("@")[1].split("/")[0].split(":")[0] + if host_part and host_part not in {"localhost", "127.0.0.1"}: + return f"http://{host_part}" + + return "http://127.0.0.1" + + +def _is_remote_worker() -> bool: + """True when this process is configured to talk to a remote main server.""" + url = _resolve_cluster_server_url() + parsed = urllib.parse.urlparse(url) + host = (parsed.hostname or "").strip().lower() + return host not in {"", "127.0.0.1", "localhost"} + + +async def materialize_cluster_input( + item: DinsarProductionRunItemORM, + source_task_dir: str, + task_id: str, +) -> None: + """Download and extract the input data for a cluster item. + + Calls ``GET /api/cluster/input-package/{item_id}`` on the main + server, retrieves a zip containing the Task_Pool directory tree, + and extracts it so that *source_task_dir* exists locally. + """ + from .task_service import task_service + + server_url = _resolve_cluster_server_url() + download_url = f"{server_url}/api/cluster/input-package/{item.id}" + parent_dir = os.path.dirname(source_task_dir) + task_name = os.path.basename(source_task_dir) + + await task_service.add_log( + task_id, + "INFO", + f"[cluster] Downloading input data from {download_url} ...", + ) + + tmp_zip = os.path.join( + tempfile.gettempdir(), + f"cluster_input_{item.id}_{task_name}.zip", + ) + try: + req = urllib.request.Request( + download_url, + headers=_cluster_request_headers(), + method="GET", + ) + with urllib.request.urlopen(req, timeout=_cluster_transfer_timeout()) as resp: + with open(tmp_zip, "wb") as fh: + shutil.copyfileobj(resp, fh, 8 * 1024 * 1024) + + os.makedirs(parent_dir, exist_ok=True) + with zipfile.ZipFile(tmp_zip, "r") as zf: + safe_extract_zip(zf, parent_dir) + + if not os.path.isdir(source_task_dir): + raise RuntimeError( + f"Extraction did not create expected directory: {source_task_dir}" + ) + + await task_service.add_log( + task_id, + "INFO", + f"[cluster] Input data ready: {source_task_dir}", + ) + except urllib.error.HTTPError as exc: + body_text = "" + try: + body_text = exc.read().decode("utf-8", errors="replace") + except Exception: + pass + await task_service.add_log( + task_id, + "WARNING", + f"[cluster] Input download HTTP {exc.code}: {body_text[:300]}", + ) + raise RuntimeError( + f"Input download failed HTTP {exc.code}: {body_text[:200]}" + ) from exc + finally: + try: + os.unlink(tmp_zip) + except Exception: + pass + + +async def upload_cluster_result( + item: DinsarProductionRunItemORM, + run: DinsarProductionRunORM, + managed_run_dir: str, + run_key: str, + task_id: str, +) -> bool: + """Package the managed run directory and upload it to the main server. + + Returns ``True`` when the main server accepted the upload and + registered the result in the D-InSAR catalog. + """ + from .task_service import task_service + + server_url = _resolve_cluster_server_url() + + await task_service.add_log( + task_id, + "INFO", + "[cluster] Packaging results for upload ...", + ) + + run_dir_name = os.path.basename(os.path.normpath(managed_run_dir)) + tmp_zip = os.path.join( + tempfile.gettempdir(), + f"cluster_result_{item.id}.zip", + ) + try: + parent = os.path.dirname(managed_run_dir) + shutil.make_archive( + tmp_zip.replace(".zip", ""), + "zip", + root_dir=parent, + base_dir=run_dir_name, + ) + + await task_service.add_log( + task_id, + "INFO", + "[cluster] Uploading results to main server ...", + ) + + boundary = "----ClusterUploadBoundary" + body = io.BytesIO() + + def _write_field(name, filename, content_type, data): + body.write(f"--{boundary}\r\n".encode("utf-8")) + if filename: + body.write( + f'Content-Disposition: form-data; name="{name}"; ' + f'filename="{filename}"\r\n'.encode("utf-8") + ) + body.write(f"Content-Type: {content_type}\r\n".encode("utf-8")) + else: + body.write( + f'Content-Disposition: form-data; name="{name}"\r\n'.encode( + "utf-8" + ) + ) + body.write(b"\r\n") + body.write(data) + body.write(b"\r\n") + + _write_field("run_id", "", "text/plain", (run.run_id or "").encode("utf-8")) + _write_field("run_key", "", "text/plain", str(run_key or "").encode("utf-8")) + with open(tmp_zip, "rb") as fh: + _write_field( + "result_zip", + os.path.basename(tmp_zip), + "application/zip", + fh.read(), + ) + body.write(f"--{boundary}--\r\n".encode("utf-8")) + + upload_url = f"{server_url}/api/cluster/upload-result/{item.id}" + req = urllib.request.Request( + upload_url, + data=body.getvalue(), + headers={ + "Content-Type": f"multipart/form-data; boundary={boundary}", + **_cluster_request_headers(), + }, + method="POST", + ) + + with urllib.request.urlopen(req, timeout=_cluster_transfer_timeout()) as resp: + result = json.loads(resp.read().decode("utf-8")) + await task_service.add_log( + task_id, + "INFO", + f"[cluster] Results uploaded: " + f"registered={result.get('registered', False)} " + f"processed={result.get('processed', 0)}", + ) + return bool(result.get("registered", False)) + + except urllib.error.HTTPError as exc: + body_text = "" + try: + body_text = exc.read().decode("utf-8", errors="replace") + except Exception: + pass + await task_service.add_log( + task_id, + "WARNING", + f"[cluster] Upload HTTP {exc.code}: {body_text[:300]}", + ) + raise RuntimeError( + f"Upload failed HTTP {exc.code}: {body_text[:200]}" + ) from exc + finally: + try: + if os.path.isfile(tmp_zip): + os.unlink(tmp_zip) + except Exception: + pass diff --git a/backend/app/services/health_service.py b/backend/app/services/health_service.py index fa2670d..617b248 100644 --- a/backend/app/services/health_service.py +++ b/backend/app/services/health_service.py @@ -20,6 +20,7 @@ from ..models import ( SARSceneGeoORM, SceneOrbitBindingORM, SourceProductAssetORM, + SystemTaskORM, SystemWorkerHeartbeatORM, ) from ..idl_service import get_idl_status @@ -1434,6 +1435,50 @@ async def _check_wsl_runtime() -> Dict[str, Any]: return status + +STUCK_TASK_THRESHOLD_SECONDS = read_int_env( + "HEALTH_STUCK_TASK_THRESHOLD_SECONDS", + 3600, + minimum=300, + maximum=86400, +) + + +async def _check_stuck_tasks() -> Dict[str, Any]: + """Detect RUNNING tasks whose updated_at has not changed for too long.""" + from datetime import timedelta + + from ..database import AsyncSessionLocal + async with AsyncSessionLocal() as db: + now = datetime.utcnow() + cutoff = now - timedelta(seconds=STUCK_TASK_THRESHOLD_SECONDS) + result = await db.execute( + select(SystemTaskORM).where( + SystemTaskORM.status == "RUNNING", + SystemTaskORM.updated_at < cutoff, + ).order_by(SystemTaskORM.updated_at.asc()) + ) + stuck = result.scalars().all() + items = [] + for task in stuck: + minutes_stuck = max(0.0, (now - task.updated_at).total_seconds()) / 60.0 + items.append({ + "task_id": task.task_id, + "task_type": task.task_type, + "task_name": task.task_name, + "progress": task.progress, + "message": task.message, + "stuck_minutes": round(minutes_stuck, 1), + "started_at": task.started_at.isoformat() if task.started_at else None, + "last_updated_at": task.updated_at.isoformat() if task.updated_at else None, + }) + return { + "ok": len(items) == 0, + "stuck_count": len(items), + "threshold_seconds": STUCK_TASK_THRESHOLD_SECONDS, + "stuck_tasks": items, + } + async def get_health_status( include_external: bool = True, include_details: bool = False, @@ -1464,6 +1509,7 @@ async def get_health_status( asset_inventory_status = await _check_asset_inventory() wsl_runtime_status = await _check_wsl_runtime() pairing_system_status = await pairing_state_service.get_pairing_system_status() + stuck_task_status = await _check_stuck_tasks() engines_status = {"ok": None, "overall": None, "engines": []} if full or include_details: engines_status = await _check_dinsar_engines() @@ -1482,6 +1528,7 @@ async def get_health_status( asset_inventory_status.get("ok"), wsl_runtime_status.get("ok"), pairing_system_status.get("ok"), + stuck_task_status.get("ok"), (not settings.TIMESERIES_ENABLED) or timeseries_result_catalog_status.get("ok"), (not (settings.GAMMA_SBAS_ENABLED or settings.LANDSAR_SBAS_ENABLED)) or sbas_insar_result_catalog_status.get("ok"), @@ -1511,6 +1558,7 @@ async def get_health_status( "asset_inventory": asset_inventory_status, "wsl_runtime": wsl_runtime_status, "pairing_system": pairing_system_status, + "stuck_tasks": stuck_task_status, "idl": { "ok": idl_ok, "status": idl_status, diff --git a/backend/app/services/job_handlers.py b/backend/app/services/job_handlers.py index cb0fc2a..6a300be 100644 --- a/backend/app/services/job_handlers.py +++ b/backend/app/services/job_handlers.py @@ -23,6 +23,11 @@ from ..models import SystemJobORM, DinsarResultORM, HazardPointORM, DinsarTaskIt from ..scheduler import scan_data_job from .data_service import data_service from .asset_inventory_service import asset_inventory_service +from .cluster_transport import ( + _is_remote_worker, + materialize_cluster_input, + upload_cluster_result, +) from .dinsar_compat_service import dinsar_compat_service from .dinsar_naming import build_run_key from .dinsar_production_service import dinsar_production_service @@ -156,6 +161,22 @@ def _normalize_positive_int(value: Any) -> Optional[int]: return parsed if parsed > 0 else None +def _cluster_source_task_dir_ready(path: str) -> bool: + if not path or not os.path.isdir(path): + return False + for child_name in ("master", "slave"): + child_dir = os.path.join(path, child_name) + if not os.path.isdir(child_dir): + return False + try: + with os.scandir(child_dir) as entries: + if not any(entries): + return False + except OSError: + return False + return True + + def _dedupe_existing_dirs(paths: Any) -> List[str]: ordered: List[str] = [] for raw_path in paths or []: @@ -3254,6 +3275,10 @@ async def _handle_landsar_cluster_item(job: SystemJobORM) -> None: raise ValueError(f"LandSAR cluster item not found: {item_id}") await db.refresh(item) + source_task_dir = os.path.normpath(str(item.source_task_dir or "")) + if source_task_dir and not _cluster_source_task_dir_ready(source_task_dir): + await materialize_cluster_input(item, source_task_dir, job.task_id) + item_status = str(item.status or "").strip().upper() if item_status in {"COMPLETED", "FAILED", "SKIPPED", "CANCELLED"}: await dinsar_production_service.finalize_cluster_run_if_complete(run, db=db) @@ -3491,26 +3516,34 @@ async def _handle_landsar_cluster_item(job: SystemJobORM) -> None: db=db, ) - try: - publish_result = await result_catalog_service.publish_from_sources(db, [managed_run_dir]) - processed_count = int(publish_result.get("processed", 0) or 0) - failed_count = int(publish_result.get("failed", 0) or 0) - if processed_count > 0: - await result_catalog_service.rebuild_catalog(db, full_rebuild=True) - if processed_count != 1 or failed_count != 0: - raise RuntimeError(f"expected processed=1 failed=0, got processed={processed_count} failed={failed_count}") - await task_service.add_log( - job.task_id, - "INFO", - f"[cluster {item_index}/{total_items}] Published {item_label}", - ) - except Exception as exc: - publish_error = str(exc) - await task_service.add_log( - job.task_id, - "WARNING", - f"[cluster {item_index}/{total_items}] Result catalog publish failed for {item_label}: {publish_error}", + # ---- Post-flight: upload or local publish ---- + if _is_remote_worker(): + upload_success = await upload_cluster_result( + item, run, managed_run_dir, run_key, job.task_id, ) + if not upload_success: + raise RuntimeError("Cluster result upload failed or not registered.") + else: + try: + publish_result = await result_catalog_service.publish_from_sources(db, [managed_run_dir]) + processed_count = int(publish_result.get("processed", 0) or 0) + failed_count = int(publish_result.get("failed", 0) or 0) + if processed_count > 0: + await result_catalog_service.rebuild_catalog(db, full_rebuild=True) + if processed_count != 1 or failed_count != 0: + raise RuntimeError(f"expected processed=1 failed=0, got processed={processed_count} failed={failed_count}") + await task_service.add_log( + job.task_id, + "INFO", + f"[cluster {item_index}/{total_items}] Published {item_label}", + ) + except Exception as exc: + publish_error = str(exc) + await task_service.add_log( + job.task_id, + "WARNING", + f"[cluster {item_index}/{total_items}] Result catalog publish failed for {item_label}: {publish_error}", + ) await task_service.add_log( job.task_id, diff --git a/config/landsar_cluster_worker.env.example b/config/landsar_cluster_worker.env.example index 7c2f091..085a407 100644 --- a/config/landsar_cluster_worker.env.example +++ b/config/landsar_cluster_worker.env.example @@ -5,3 +5,24 @@ JOB_WORKER_ALLOWED_TYPES=LANDSAR_CLUSTER_ITEM JOB_WORKER_CONCURRENCY=1 JOB_WORKER_POLL_INTERVAL=1.0 LANDSAR_CLUSTER_WORKER_ID= +CLUSTER_MAIN_SERVER_URL=http://192.168.1.62 +CLUSTER_SHARED_TOKEN= +CLUSTER_TRANSFER_TIMEOUT_SECONDS=3600 + +PYTHON_PATH=C:\ProgramData\anaconda3\envs\InSAR\python.exe +LANDSAR_ENABLED=true +LANDSAR_HOME=E:\LandSAR +LANDSAR_EXTRA_HOME=D:\Code\Insar_management_system_v2\third_party\LandSAR +LANDSAR_CONSOLE_EXE=E:\LandSAR\InSAR_Console.exe +LANDSAR_WORK_ROOT=E:\LandSAR_Work +LANDSAR_DEM_PATH=E:\DEM\SRTMDEM_RSP_SARscape_global_int16.tif +LANDSAR_LICENSE_MODE=netVersion +LANDSAR_LICENSE_HOST=127.0.0.1 +LANDSAR_LICENSE_PORT=6666 +LANDSAR_CONFIG_ROW=netVersion,zh,127.0.0.1,6666 +LANDSAR_CONFIG_AUTO_WRITE=true +LANDSAR_AUTH_SERVER_EXE=D:\Code\Insar_management_system_v2\third_party\LandSAR\tools\_portable_release\LandSAR_auth_tools_win64\landsar_net_auth_server.exe +LANDSAR_AUTH_SERVER_AUTO_START=true +LANDSAR_AUTH_SERVER_HOST=127.0.0.1 +LANDSAR_AUTH_SERVER_PORT=6666 +LANDSAR_DINSAR_TIMEOUT_SECONDS=43200 diff --git a/docs/INDEX.md b/docs/INDEX.md index efac07c..cf0d193 100644 --- a/docs/INDEX.md +++ b/docs/INDEX.md @@ -37,6 +37,8 @@ LandSAR D-InSAR/SBAS 鐨勫叏鐞?DEM 涓€娆℃€?Int16 鏍囧噯鍖栥€佸尯鍩熻鍓?tif銆佺敓浜ч厤缃拰 guardrail 绾﹀畾銆? - [LANDSAR_CLUSTER_WORKER_DEPLOYMENT_20260624.md](LANDSAR_CLUSTER_WORKER_DEPLOYMENT_20260624.md) LandSAR D-InSAR 集群 worker 的队列分片设计、主服务器 IP 白名单、远端 Windows 节点 192.168.1.6 部署和运行约束。 +- [LANDSAR_CLUSTER_DATA_TRANSPORT_DESIGN_20260625.md](LANDSAR_CLUSTER_DATA_TRANSPORT_DESIGN_20260625.md) + LandSAR 集群数据搬运(HTTP Task_Pool 下载 + 结果回传)、Windows 集群运维(Task Scheduler 开机自启 + 心跳监控)。 - [UNC_SOURCE_ARCHIVE_AND_MATERIALIZE_DESIGN_20260615.md](UNC_SOURCE_ARCHIVE_AND_MATERIALIZE_DESIGN_20260615.md) LT-1/Sentinel-1 鏈湴婧愬帇缂╁寘绠$悊銆佸寘鍐?XML/manifest 璧勪骇鍖栥€佹湰鍦?Task_Pool materialize锛屼互鍙?UNC 閫€鍑哄悗鐨勬湰鏈洪儴缃茶竟鐣屻€? - [SOURCE_ARCHIVE_INTEGRITY_AUDIT_20260620.md](SOURCE_ARCHIVE_INTEGRITY_AUDIT_20260620.md) diff --git a/docs/LANDSAR_CLUSTER_DATA_TRANSPORT_DESIGN_20260625.md b/docs/LANDSAR_CLUSTER_DATA_TRANSPORT_DESIGN_20260625.md new file mode 100644 index 0000000..c7444cf --- /dev/null +++ b/docs/LANDSAR_CLUSTER_DATA_TRANSPORT_DESIGN_20260625.md @@ -0,0 +1,235 @@ + # LandSAR 集群数据搬运与运维设计(2026-06-25) + + ## 背景 + + [LANDSAR_CLUSTER_WORKER_DEPLOYMENT_20260624.md](LANDSAR_CLUSTER_WORKER_DEPLOYMENT_20260624.md) 已完成集群调度骨架:主服务器提交 LandSAR 集群任务 → 按 pair 拆分为 `LANDSAR_CLUSTER_ITEM` → 本机或远端 worker 通过 DB 队列领取执行。 + + 该文档留下了两个明确缺口: + + 1. **数据搬运**:远端 worker 执行 `engine.run()` 时读取 `item.source_task_dir`(如 `D:\Task_Pool\DInSAR\Task_20250601_20250612\master`),该路径在远端不存在。 + 2. **结果回传**:远端 worker 处理完成后,标准产品包在远端本地磁盘,不会自动进入主服务器 D-InSAR catalog。 + + 本文档定义这两个能力的设计,以及 Windows 集群运维方案。 + + ## 设计目标 + + - 不依赖 Windows 文件共享(SMB)、映射盘符、UNC 路径 + - 不要求主服务器和 worker 共用盘符或路径结构 + - Worker 节点可以动态增减,配置简单 + - 复用现有 `SOURCE_PRODUCT_DIRS` 压缩包源池和 `TASK_POOL_ROOT` 体系 + - 传输失败利用队列系统自带的重试机制 + + ## 架构总览 + + ``` + 主服务器 (192.168.1.62) 远端 Worker (192.168.1.6 / .7 / ...) + ════════════════════════ ═══════════════════════════════════ + + [前端] 提交集群任务 + │ + ▼ + [Router] POST /landsar-cluster/run + │ + ▼ + [Service] create_landsar_cluster_run() + │ 扫描 Task_* → 为每个 pair 创建 + │ DinsarProductionRunItemORM + + │ LANDSAR_CLUSTER_ITEM 队列任务 + ▼ + [system_jobs 表] ◄─────────── claim_next_job ────── [run_landsar_cluster_worker.py] + │ + ┌────▼──────────────────┐ + │ 检查 source_task_dir │ + │ 本地是否存在 │ + ├────有─────────────────┤ + │ → 跳到 LandSAR 执行 │ + ├────无─────────────────┤ + │ 1. GET /api/cluster/ │ + │ input-package/ │ + │ 2. 解包到本地路径 │ + └───────────────────────┘ + │ + ▼ + LandSAR engine.run() + │ + ▼ + ┌───────────────────────┐ + │ POST /api/cluster/ │ + │ upload-result/ │ + │ → 主服务器写入结果目录 │ + │ → catalog 登记 │ + └───────────────────────┘ + ``` + + ## 数据搬运详细设计 + + ### 阶段 1: Pre-flight 输入数据下载 + + **触发时机**:`_handle_landsar_cluster_item` 执行 `engine.run()` 之前。 + + **Worker 端流程**: + + 1. 读取 `item.source_task_dir`,检查本地目录是否存在且包含 `master/`、`slave/`、`pair_metadata.json` + 2. 若无 → 调用 `GET /api/cluster/input-package/{item_id}` + 3. Worker 请求头携带 `X-Cluster-Token`,主服务器校验 `CLUSTER_SHARED_TOKEN` + 4. 主服务器读取 `item.source_task_dir` 指向的现有 Task_Pool 目录 + 5. 将该 Task_Pool 目录打成 zip 流式返回给 worker + 6. Worker 解包到 `item.source_task_dir`(保持与主服务器一致的路径结构) + 7. Worker 对 zip 成员路径做目录逃逸校验,并校验解包后的 `master/`、`slave/` 目录可用 + + **主服务器新增 API**: + + ``` + GET /api/cluster/input-package/{item_id} + Header: X-Cluster-Token: + Response: application/zip (streaming) + 内容结构: + Task_YYYYMMDD_YYYYMMDD/ + master/ + <源文件...> + slave/ + <源文件...> + orbit/ + <精轨文件...> + pair_metadata.json + ``` + + **当前实现边界**:主服务器不在传输接口里重新从 `SOURCE_PRODUCT_DIRS` 解包源压缩包;传输接口只打包已经由 Task_Pool 准备流程生成的 `item.source_task_dir`。如果主服务器上的 `source_task_dir` 不存在,接口返回 404,队列重试会保留失败信息。 + + ### 阶段 2: Post-flight 结果上传 + + **触发时机**:LandSAR 执行成功、execution manifest 构建完成后。 + + **Worker 端流程**: + + 1. LandSAR 完成后,收集标准产品包文件列表(从 `task_result.source_files` 和 manifest 获取) + 2. 调用 `POST /api/cluster/upload-result/{item_id}` + 3. 以 multipart 或流式上传产品包(含 primary file、auxiliary files、metadata) + 4. 主服务器接收后写入该 item 的标准目录 `results_root_dir\runs\` + 5. 触发 catalog 登记(复用现有 `result_catalog_service.bootstrap` 或增量登记) + + **主服务器新增 API**: + + ``` + POST /api/cluster/upload-result/{item_id} + Header: X-Cluster-Token: + Body: multipart/form-data + - result_zip: managed run directory zip + - run_id: 集群 run ID + - run_key: 当前执行 run key + Response: { "registered": true, "processed": 1, "failed": 0, "catalog_path": "..." } + ``` + + **上传后处理**: + + 1. 校验文件完整性(与 manifest 对比) + 2. 写入 `D:\production_results\dinsar\\\\` + 3. 调用 `result_catalog_service` 增量登记 + 4. 更新 `DinsarProductionExecutionORM` 指向最终结果路径 + 5. 返回登记结果给 worker + + ### 错误处理与重试 + + - 下载失败:worker 端抛异常 → `job_queue_service.mark_failed` → 按 `max_attempts` 自动重试 + - 上传失败:同上,LandSAR 已完成的中间结果在 worker 本地保留(下次重试跳过 LandSAR,直接上传) + - 主服务器端打包失败:返回 500 + 错误详情,worker 捕获后走重试 + - 超时控制:复用 `LANDSAR_DINSAR_TIMEOUT_SECONDS`,传输阶段额外设 `CLUSTER_TRANSFER_TIMEOUT_SECONDS`(默认 3600) + + ## Windows 集群运维设计 + + ### Worker 开机自启 + + 每个 worker 节点配置一条 **Windows Task Scheduler** 任务: + + ```powershell + # 创建计划任务(以管理员身份运行) + $action = New-ScheduledTaskAction -Execute "powershell.exe" ` + -Argument "-NoProfile -ExecutionPolicy Bypass -File D:\Code\Insar_management_system_v2\scripts\start_landsar_cluster_worker.ps1 -Background" + $trigger = New-ScheduledTaskTrigger -AtStartup + $settings = New-ScheduledTaskSettingsSet -RestartCount 3 -RestartInterval (New-TimeSpan -Minutes 1) ` + -AllowStartIfOnBatteries -DontStopIfGoingOnBatteries + Register-ScheduledTask -TaskName "InSAR_LandSAR_Cluster_Worker" ` + -Action $action -Trigger $trigger -Settings $settings ` + -RunLevel Highest -Description "LandSAR 集群 Worker 常驻进程" + ``` + + 关键设置: + - 触发器:系统启动时 + - 失败重试:3 次,间隔 1 分钟 + - 运行级别:最高权限 + - 不因电池模式停止(针对笔记本) + + ### Worker 监控 + + Worker 进程已内建心跳机制(`_touch_worker`),每隔 `JOB_WORKER_HEARTBEAT_INTERVAL` 秒向 `system_worker_heartbeats` 表写入心跳。主服务器健康检查页面可以展示: + + - 各 worker 的 hostname / PID / worker_id + - 最后心跳时间 + - 当前正在处理的 job 数量 + - 历史完成/失败统计 + + ### 代码同步 + + 远端 worker 建议通过 Git 同步代码: + + ```powershell + # 在远端 192.168.1.6 上 + cd D:\Code\Insar_management_system_v2 + git fetch origin + git checkout + ``` + + `.env` 文件不在 Git 中,需手动维护。远端 `.env` 的 `DATABASE_URL` 指向主服务器,LandSAR 路径指向远端本地。 + + ## 配置新增项 + + 主服务器 `.env` 新增: + + ```env + # 集群传输配置 + CLUSTER_SHARED_TOKEN= + CLUSTER_TRANSFER_TIMEOUT_SECONDS=3600 + CLUSTER_MATERIALIZE_TEMP_DIR=D:\Task_Pool\_cluster_temp + ``` + + Worker 端 `.env` 新增(或保持现有模板字段): + + ```env + # 集群 worker 主服务器地址 + CLUSTER_MAIN_SERVER_URL=http://192.168.1.62 + CLUSTER_SHARED_TOKEN= + CLUSTER_TRANSFER_TIMEOUT_SECONDS=3600 + ``` + + `CLUSTER_MAIN_SERVER_URL` 用于 worker 构造下载/上传 API 的完整 URL。未配置时默认使用 `DATABASE_URL` 中的 host 推断。 + `CLUSTER_SHARED_TOKEN` 是集群传输接口的专用共享密钥,主服务器和所有远端 worker 必须一致且非空;未配置时 `/api/cluster/...` 接口返回 503。 + + ## 实现路线图 + + | 阶段 | 内容 | 依赖 | + | --- | --- | --- | + | 1 | 主服务器 `GET /api/cluster/input-package/{item_id}` | Task_Pool `source_task_dir` | + | 2 | Worker handler 增加 pre-flight download + extract | 阶段 1 | + | 3 | 主服务器 `POST /api/cluster/upload-result/{item_id}` | `result_catalog_service` | + | 4 | Worker handler 增加 post-flight upload | 阶段 3 | + | 5 | 端到端测试(主服务器 + 远端 .6) | 阶段 1-4 | + | 6 | Windows Task Scheduler 开机自启配置 | 阶段 5 | + + ## 当前状态(2026-06-26) + + - [x] 集群调度骨架(提交 → 拆 item → DB 队列 → worker 领取 → LandSAR 执行) + - [x] Worker 入口脚本 + 启动/停止脚本 + - [x] Worker 心跳上报 + - [x] 本机集群模式验证通过(`SameSite=lax` Cookie 修复后) + - [x] 数据搬运 pre-flight(HTTP 下载 Task_Pool zip + 安全解压) + - [x] 结果上传 post-flight(HTTP 上传结果 zip + catalog 登记) + - [x] 集群传输接口 `CLUSTER_SHARED_TOKEN` 鉴权 + - [ ] 远端 .6 端到端测试 + - [ ] Worker 开机自启 + + ## 相关文档 + + - [LANDSAR_CLUSTER_WORKER_DEPLOYMENT_20260624.md](LANDSAR_CLUSTER_WORKER_DEPLOYMENT_20260624.md) — 集群调度骨架与部署记录 + - [LANDSAR_DEM_PREPARATION_CONTRACT_20260618.md](LANDSAR_DEM_PREPARATION_CONTRACT_20260618.md) — LandSAR DEM 准备约定 + - [THREE_SENSOR_LOCAL_PRODUCTION_CONTRACT_20260616.md](THREE_SENSOR_LOCAL_PRODUCTION_CONTRACT_20260616.md) — 三数据本机生产约定 + - [DINSAR_TASK_POOL_THREE_ENGINE_REFACTOR_20260614.md](DINSAR_TASK_POOL_THREE_ENGINE_REFACTOR_20260614.md) — Task_Pool 与三引擎 refactor diff --git a/scripts/start_landsar_cluster_worker.ps1 b/scripts/start_landsar_cluster_worker.ps1 index 8eb4e84..cf1b626 100644 --- a/scripts/start_landsar_cluster_worker.ps1 +++ b/scripts/start_landsar_cluster_worker.ps1 @@ -52,6 +52,162 @@ function Resolve-PythonPath { throw "Python interpreter not found. Set PYTHON_PATH in .env or pass -PythonPath." } +function Read-BoolDotEnvValue { + param( + [string]$Path, + [string]$Name, + [bool]$DefaultValue + ) + $value = Read-DotEnvValue -Path $Path -Name $Name + if (-not $value) { + return $DefaultValue + } + return @("1", "true", "yes", "on") -contains $value.Trim().ToLowerInvariant() +} + +function Read-IntDotEnvValue { + param( + [string]$Path, + [string]$Name, + [int]$DefaultValue + ) + $value = Read-DotEnvValue -Path $Path -Name $Name + if (-not $value) { + return $DefaultValue + } + $parsed = 0 + if ([int]::TryParse($value.Trim(), [ref]$parsed)) { + return $parsed + } + return $DefaultValue +} + +function Test-TcpPort { + param( + [string]$HostName, + [int]$Port + ) + try { + $client = [System.Net.Sockets.TcpClient]::new() + $async = $client.BeginConnect($HostName, $Port, $null, $null) + $ok = $async.AsyncWaitHandle.WaitOne(1000, $false) + if ($ok) { + $client.EndConnect($async) + } + $client.Close() + return $ok + } catch { + return $false + } +} + +function Get-LandSARConfigEndpoint { + param( + [string]$Path + ) + $row = Read-DotEnvValue -Path $Path -Name "LANDSAR_CONFIG_ROW" + if ($row) { + $parts = $row.Split(",") | ForEach-Object { $_.Trim() } + if ($parts.Count -ge 4) { + $port = 6666 + [void][int]::TryParse($parts[3], [ref]$port) + return [pscustomobject]@{ + Mode = $parts[0] + Host = $(if ($parts[2]) { $parts[2] } else { "127.0.0.1" }) + Port = $port + } + } + } + return [pscustomobject]@{ + Mode = $(Read-DotEnvValue -Path $Path -Name "LANDSAR_LICENSE_MODE") + Host = $(Read-DotEnvValue -Path $Path -Name "LANDSAR_LICENSE_HOST") + Port = $(Read-IntDotEnvValue -Path $Path -Name "LANDSAR_LICENSE_PORT" -DefaultValue 6666) + } +} + +function Resolve-LandSARAuthServerPath { + param( + [string]$RepoRoot, + [string]$EnvPath + ) + $explicit = Read-DotEnvValue -Path $EnvPath -Name "LANDSAR_AUTH_SERVER_EXE" + $candidates = @( + $explicit, + (Join-Path $RepoRoot "third_party\LandSAR\tools\_portable_release\LandSAR_auth_tools_win64\landsar_net_auth_server.exe"), + (Join-Path $RepoRoot "third_party\LandSAR\landsar_net_auth_server.exe") + ) | Where-Object { $_ -and $_.Trim() } | Select-Object -Unique + foreach ($candidate in $candidates) { + if (Test-Path -LiteralPath $candidate) { + return (Resolve-Path -LiteralPath $candidate).Path + } + } + return $candidates | Select-Object -First 1 +} + +function Start-LandSARAuthServerIfNeeded { + param( + [string]$RepoRoot, + [string]$EnvPath + ) + $endpoint = Get-LandSARConfigEndpoint -Path $EnvPath + $mode = $(if ($endpoint.Mode) { $endpoint.Mode } else { "netVersion" }) + if ($mode.Trim().ToLowerInvariant() -ne "netversion") { + Write-Host "LandSAR auth: skipped for license mode $mode" + return + } + + $clientHost = $(if ($endpoint.Host) { $endpoint.Host } else { "127.0.0.1" }) + $clientPort = [int]$endpoint.Port + if (Test-TcpPort -HostName $clientHost -Port $clientPort) { + Write-Host "LandSAR auth: already listening on $clientHost`:$clientPort" + return + } + + $autoStart = Read-BoolDotEnvValue -Path $EnvPath -Name "LANDSAR_AUTH_SERVER_AUTO_START" -DefaultValue $true + if (-not $autoStart) { + throw "LandSAR auth server is not listening on $clientHost`:$clientPort and LANDSAR_AUTH_SERVER_AUTO_START=false." + } + + $authExe = Resolve-LandSARAuthServerPath -RepoRoot $RepoRoot -EnvPath $EnvPath + if (-not $authExe -or -not (Test-Path -LiteralPath $authExe)) { + throw "LandSAR auth server executable not found: $authExe. Set LANDSAR_AUTH_SERVER_EXE in .env." + } + + $serverDir = Split-Path -Parent $authExe + $memoryBin = Join-Path $serverDir "dongle_0xa0.bin" + if (-not (Test-Path -LiteralPath $memoryBin)) { + $fallbackBin = Join-Path $RepoRoot "third_party\LandSAR\tools\dongle_0xa0.bin" + if (Test-Path -LiteralPath $fallbackBin) { + Copy-Item -LiteralPath $fallbackBin -Destination $memoryBin -Force + } else { + throw "LandSAR auth memory image missing: $memoryBin" + } + } + + $bindHost = Read-DotEnvValue -Path $EnvPath -Name "LANDSAR_AUTH_SERVER_HOST" + if (-not $bindHost) { + $bindHost = $clientHost + } + $bindPort = Read-IntDotEnvValue -Path $EnvPath -Name "LANDSAR_AUTH_SERVER_PORT" -DefaultValue $clientPort + + Write-Host "LandSAR auth: starting $authExe on $bindHost`:$bindPort" + Start-Process ` + -FilePath $authExe ` + -ArgumentList @("--host", $bindHost, "--port", [string]$bindPort) ` + -WorkingDirectory $serverDir ` + -WindowStyle Hidden | Out-Null + + $deadline = (Get-Date).AddSeconds(5) + while ((Get-Date) -lt $deadline) { + if (Test-TcpPort -HostName $clientHost -Port $clientPort) { + Write-Host "LandSAR auth: started and reachable on $clientHost`:$clientPort" + return + } + Start-Sleep -Milliseconds 250 + } + throw "LandSAR auth server was started but $clientHost`:$clientPort is not reachable." +} + $RepoRoot = (Resolve-Path -LiteralPath $RepoRoot).Path $envPath = Join-Path $RepoRoot ".env" $templatePath = Join-Path $RepoRoot "config\landsar_cluster_worker.env.example" @@ -76,6 +232,10 @@ $allowedTypes = Read-DotEnvValue -Path $envPath -Name "JOB_WORKER_ALLOWED_TYPES" if (-not $allowedTypes) { $env:JOB_WORKER_ALLOWED_TYPES = "LANDSAR_CLUSTER_ITEM" } +$clusterToken = Read-DotEnvValue -Path $envPath -Name "CLUSTER_SHARED_TOKEN" +if (-not $clusterToken) { + throw "CLUSTER_SHARED_TOKEN is required for LandSAR cluster input download and result upload." +} $concurrency = Read-DotEnvValue -Path $envPath -Name "JOB_WORKER_CONCURRENCY" if (-not $concurrency) { $env:JOB_WORKER_CONCURRENCY = "1" @@ -101,6 +261,7 @@ Write-Host "Log: $stdoutLog" Write-Host "Mode: $(if ($Background) { 'background' } else { 'foreground' })" Set-Location $RepoRoot +Start-LandSARAuthServerIfNeeded -RepoRoot $RepoRoot -EnvPath $envPath if ($Background) { if (Test-Path -LiteralPath $pidFile) { diff --git a/scripts/test_landsar_cluster_worker_remote.ps1 b/scripts/test_landsar_cluster_worker_remote.ps1 new file mode 100644 index 0000000..a300f17 --- /dev/null +++ b/scripts/test_landsar_cluster_worker_remote.ps1 @@ -0,0 +1,157 @@ +param( + [string]$RepoRoot = (Resolve-Path (Join-Path $PSScriptRoot "..")).Path +) + +$ErrorActionPreference = "Stop" + +function Read-DotEnvValue { + param( + [string]$Path, + [string]$Name + ) + if (-not (Test-Path -LiteralPath $Path)) { + return "" + } + $line = Get-Content -LiteralPath $Path | + Where-Object { $_ -match "^\s*$([regex]::Escape($Name))\s*=" } | + Select-Object -Last 1 + if (-not $line) { + return "" + } + $value = ($line -split "=", 2)[1].Trim() + if (($value.StartsWith('"') -and $value.EndsWith('"')) -or ($value.StartsWith("'") -and $value.EndsWith("'"))) { + $value = $value.Substring(1, $value.Length - 2) + } + return $value +} + +function Add-Check { + param( + [System.Collections.Generic.List[object]]$Rows, + [string]$Name, + [bool]$Ok, + [string]$Detail + ) + $Rows.Add([pscustomobject]@{ + Check = $Name + Ok = $Ok + Detail = $Detail + }) | Out-Null +} + +$RepoRoot = (Resolve-Path -LiteralPath $RepoRoot).Path +$envPath = Join-Path $RepoRoot ".env" +$rows = [System.Collections.Generic.List[object]]::new() + +Add-Check $rows "repo_root" (Test-Path -LiteralPath $RepoRoot) $RepoRoot +Add-Check $rows "env_file" (Test-Path -LiteralPath $envPath) $envPath +Add-Check $rows "worker_script" (Test-Path -LiteralPath (Join-Path $RepoRoot "run_landsar_cluster_worker.py")) (Join-Path $RepoRoot "run_landsar_cluster_worker.py") +Add-Check $rows "start_launcher" (Test-Path -LiteralPath (Join-Path $RepoRoot "scripts\start_landsar_cluster_worker.ps1")) (Join-Path $RepoRoot "scripts\start_landsar_cluster_worker.ps1") +Add-Check $rows "stop_launcher" (Test-Path -LiteralPath (Join-Path $RepoRoot "scripts\stop_landsar_cluster_worker.ps1")) (Join-Path $RepoRoot "scripts\stop_landsar_cluster_worker.ps1") + +$databaseUrl = Read-DotEnvValue -Path $envPath -Name "DATABASE_URL" +$dbHost = "" +$dbPort = 5432 +if ($databaseUrl -match "@(?[^:/]+)(:(?\d+))?/") { + $dbHost = $Matches.host + if ($Matches.port) { + $dbPort = [int]$Matches.port + } +} +Add-Check $rows "database_url" ([bool]$databaseUrl) $databaseUrl +if ($dbHost) { + $tcpOk = $false + try { + $tcpOk = [bool](Test-NetConnection -ComputerName $dbHost -Port $dbPort -InformationLevel Quiet -WarningAction SilentlyContinue) + } catch { + $tcpOk = $false + } + Add-Check $rows "database_tcp" $tcpOk "$dbHost`:$dbPort" +} else { + Add-Check $rows "database_tcp" $false "DATABASE_URL host could not be parsed" +} + +$pythonPath = Read-DotEnvValue -Path $envPath -Name "PYTHON_PATH" +if (-not $pythonPath) { + $pythonPath = "C:\ProgramData\anaconda3\envs\InSAR\python.exe" +} +Add-Check $rows "python_path" (Test-Path -LiteralPath $pythonPath) $pythonPath + +$landsarConsole = Read-DotEnvValue -Path $envPath -Name "LANDSAR_CONSOLE_EXE" +if (-not $landsarConsole) { + $landsarConsole = "D:\LandSAR\InSAR_Console.exe" +} +Add-Check $rows "landsar_console" (Test-Path -LiteralPath $landsarConsole) $landsarConsole + +$landsarHome = Read-DotEnvValue -Path $envPath -Name "LANDSAR_HOME" +if (-not $landsarHome) { + $landsarHome = Split-Path -Parent $landsarConsole +} +Add-Check $rows "landsar_home" (Test-Path -LiteralPath $landsarHome) $landsarHome + +$landsarExtraHome = Read-DotEnvValue -Path $envPath -Name "LANDSAR_EXTRA_HOME" +if (-not $landsarExtraHome) { + $landsarExtraHome = Join-Path $RepoRoot "third_party\LandSAR" +} +Add-Check $rows "landsar_extra_home" (Test-Path -LiteralPath $landsarExtraHome) $landsarExtraHome + +$authExe = Read-DotEnvValue -Path $envPath -Name "LANDSAR_AUTH_SERVER_EXE" +if (-not $authExe) { + $authExe = Join-Path $RepoRoot "third_party\LandSAR\tools\_portable_release\LandSAR_auth_tools_win64\landsar_net_auth_server.exe" +} +Add-Check $rows "landsar_auth_exe" (Test-Path -LiteralPath $authExe) $authExe + +$authMemory = Join-Path (Split-Path -Parent $authExe) "dongle_0xa0.bin" +$fallbackMemory = Join-Path $RepoRoot "third_party\LandSAR\tools\dongle_0xa0.bin" +Add-Check $rows "landsar_auth_memory" ((Test-Path -LiteralPath $authMemory) -or (Test-Path -LiteralPath $fallbackMemory)) "$authMemory or $fallbackMemory" + +$authHost = Read-DotEnvValue -Path $envPath -Name "LANDSAR_AUTH_SERVER_HOST" +if (-not $authHost) { + $authHost = "127.0.0.1" +} +$authPortText = Read-DotEnvValue -Path $envPath -Name "LANDSAR_AUTH_SERVER_PORT" +$authPort = 6666 +if ($authPortText) { + [void][int]::TryParse($authPortText, [ref]$authPort) +} +Add-Check $rows "landsar_auth_endpoint" ($authPort -gt 0) "$authHost`:$authPort" + +$demPath = Read-DotEnvValue -Path $envPath -Name "LANDSAR_DEM_PATH" +if (-not $demPath) { + $demPath = "D:\DEM\SRTMDEM_RSP_SARscape_global_int16.tif" +} +Add-Check $rows "landsar_dem" (Test-Path -LiteralPath $demPath) $demPath + +$workRoot = Read-DotEnvValue -Path $envPath -Name "LANDSAR_WORK_ROOT" +if (-not $workRoot) { + $workRoot = "D:\LandSAR_Work" +} +try { + New-Item -ItemType Directory -Force -Path $workRoot | Out-Null + $workRootOk = Test-Path -LiteralPath $workRoot +} catch { + $workRootOk = $false +} +Add-Check $rows "landsar_work_root" $workRootOk $workRoot + +$allowedTypes = Read-DotEnvValue -Path $envPath -Name "JOB_WORKER_ALLOWED_TYPES" +Add-Check $rows "allowed_job_types" ($allowedTypes -eq "LANDSAR_CLUSTER_ITEM") $allowedTypes + +$clusterMainUrl = Read-DotEnvValue -Path $envPath -Name "CLUSTER_MAIN_SERVER_URL" +Add-Check $rows "cluster_main_server_url" ([bool]$clusterMainUrl) $clusterMainUrl + +$clusterToken = Read-DotEnvValue -Path $envPath -Name "CLUSTER_SHARED_TOKEN" +Add-Check $rows "cluster_shared_token" ([bool]$clusterToken) $(if ($clusterToken) { "configured" } else { "missing" }) + +$rows | Format-Table -AutoSize + +$failed = @($rows | Where-Object { -not $_.Ok }) +if ($failed.Count -gt 0) { + Write-Host "" + Write-Host "FAILED CHECKS:" -ForegroundColor Red + $failed | Format-Table -AutoSize + exit 1 +} + +Write-Host "" +Write-Host "All LandSAR cluster worker checks passed." -ForegroundColor Green