Files

538 lines
22 KiB
Python

from __future__ import annotations
import os
import shutil
from datetime import datetime, timedelta, timezone
from pathlib import Path
from typing import Any, Dict, List, Optional
from sqlalchemy import func, select
from sqlalchemy.ext.asyncio import AsyncSession
from ..config import settings
from ..models import (
DinsarProductionExecutionORM,
DinsarProductionRunItemORM,
DinsarProductionRunORM,
SystemJobORM,
SystemTaskORM,
TaskLogORM,
)
MAINTENANCE_LIST_STATUSES = {"COMPLETED", "FAILED", "PARTIAL_SUCCESS", "CANCELLED", "PENDING", "RUNNING"}
ACTIVE_TASK_STATUSES = {"PENDING", "RUNNING"}
ACTIVE_JOB_STATUSES = {"READY", "RETRY", "RUNNING"}
DINSAR_TASK_TYPES = {"LANDSAR_RUN", "LANDSAR_CLUSTER_RUN", "PYINT_RUN", "IDL_RUN_DINSAR"}
SUPPORTED_CLEANUP_TASK_TYPES = {"LANDSAR_RUN", "LANDSAR_CLUSTER_RUN"}
DEFAULT_TASK_TYPES = DINSAR_TASK_TYPES | {"COPY_DATA", "PAIRING_CACHE_REBUILD"}
LANDSAR_WORK_TERMINAL_EXECUTION_STATUSES = {"COMPLETED", "FAILED", "CANCELLED"}
def _utcnow_naive() -> datetime:
return datetime.now(timezone.utc).replace(tzinfo=None)
def _dt(value: Any) -> Optional[str]:
if value is None:
return None
if hasattr(value, "isoformat"):
return value.isoformat()
return str(value)
def _norm_status(value: Any) -> str:
return str(value or "").strip().upper()
def _compact(value: Any, limit: int = 220) -> Optional[str]:
if value is None:
return None
text = str(value).replace("\r", " ").replace("\n", " ").strip()
if len(text) <= limit:
return text
return text[:limit] + "..."
def _normalize_path_text(value: Any) -> str:
return os.path.normpath(str(value or "").strip().strip('"').strip("'"))
def _path_exists(path: str) -> bool:
return bool(path) and os.path.exists(path)
def _is_within_path(path: str, parent: str) -> bool:
path_text = _normalize_path_text(path)
parent_text = _normalize_path_text(parent)
if not path_text or not parent_text:
return False
try:
return os.path.commonpath([os.path.abspath(path_text), os.path.abspath(parent_text)]) == os.path.abspath(parent_text)
except ValueError:
return False
class OpsMaintenanceService:
async def list_tasks(
self,
db: AsyncSession,
*,
task_type: Optional[str] = None,
status: Optional[str] = None,
limit: int = 50,
offset: int = 0,
) -> Dict[str, Any]:
safe_limit = max(1, min(int(limit or 50), 200))
safe_offset = max(0, int(offset or 0))
task_types = [task_type.strip().upper()] if task_type else sorted(DEFAULT_TASK_TYPES)
status_filter = status.strip().upper() if status else ""
conditions = [SystemTaskORM.task_type.in_(task_types)]
if status_filter:
conditions.append(SystemTaskORM.status == status_filter)
else:
conditions.append(SystemTaskORM.status.in_(sorted(MAINTENANCE_LIST_STATUSES)))
result = await db.execute(
select(SystemTaskORM)
.where(*conditions)
.order_by(SystemTaskORM.updated_at.desc(), SystemTaskORM.id.desc())
.offset(safe_offset)
.limit(safe_limit)
)
tasks = list(result.scalars().all())
items = [await self._build_task_summary(db, task) for task in tasks]
abnormal_items = [
item
for item in items
if item.get("issue_level") in {"warning", "danger"}
or _norm_status(item.get("status")) in {"FAILED", "PARTIAL_SUCCESS", "CANCELLED"}
]
visible_items = items if status_filter else abnormal_items
return {
"items": visible_items,
"limit": safe_limit,
"offset": safe_offset,
"returned": len(visible_items),
}
async def diagnose_task(self, db: AsyncSession, task_id: str) -> Optional[Dict[str, Any]]:
task = await self._get_task(db, task_id)
if task is None:
return None
jobs = await self._get_jobs(db, task.task_id)
run = await self._get_production_run_for_task(db, task, jobs)
item_counts: Dict[str, int] = {}
execution_counts: Dict[str, int] = {}
disk_paths: List[Dict[str, Any]] = []
if run is not None:
item_counts = await self._status_counts(db, DinsarProductionRunItemORM, run.run_id)
execution_counts = await self._status_counts(db, DinsarProductionExecutionORM, run.run_id)
disk_paths.extend(await self._collect_run_disk_paths(db, run))
recent_logs = await self._recent_logs(db, task.task_id)
findings, cleanup_supported, cleanup_blockers = self._diagnose_findings(
task=task,
jobs=jobs,
run=run,
item_counts=item_counts,
execution_counts=execution_counts,
disk_paths=disk_paths,
)
return {
"task": self._task_payload(task),
"jobs": [self._job_payload(job) for job in jobs],
"production_run": self._run_payload(run) if run else None,
"production_item_counts": item_counts,
"production_execution_counts": execution_counts,
"recent_logs": recent_logs,
"disk_paths": disk_paths,
"diagnosis": {
"summary": findings[0] if findings else "未发现明显异常。",
"findings": findings,
"cleanup_supported": cleanup_supported,
"cleanup_blockers": cleanup_blockers,
},
}
async def cleanup_preview(self, db: AsyncSession, task_id: str) -> Optional[Dict[str, Any]]:
diagnosis = await self.diagnose_task(db, task_id)
if diagnosis is None:
return None
task = diagnosis["task"]
task_type = _norm_status(task.get("task_type"))
task_status = _norm_status(task.get("status"))
cleanup_supported = bool(diagnosis["diagnosis"].get("cleanup_supported"))
blockers = list(diagnosis["diagnosis"].get("cleanup_blockers") or [])
if task_type not in SUPPORTED_CLEANUP_TASK_TYPES:
blockers.append(f"第一版暂不支持清理任务类型 {task_type}。")
if task_status in ACTIVE_TASK_STATUSES:
blockers.append("任务仍处于活动状态,不能清理。")
disk_deletes = self._cleanup_disk_targets(diagnosis)
has_existing_target = any(item.get("exists") and item.get("allowed") for item in disk_deletes)
if not has_existing_target:
blockers.append("未发现可清理的 LandSAR_WORK_ROOT/run_* 目录。")
blocked = bool(blockers) or not cleanup_supported
return {
"task_id": task_id,
"blocked": blocked,
"blockers": blockers,
"cleanup_supported": cleanup_supported and not blocked,
"database_deletes": {},
"disk_deletes": disk_deletes,
}
async def cleanup_task(
self,
db: AsyncSession,
task_id: str,
*,
options: Dict[str, bool],
) -> Optional[Dict[str, Any]]:
preview = await self.cleanup_preview(db, task_id)
if preview is None:
return None
if preview.get("blocked"):
raise ValueError("; ".join(preview.get("blockers") or ["清理被阻止。"]))
task = await self._get_task(db, task_id)
if task is None:
return None
disk_result = await self._delete_disk_targets(preview, options)
return {
"task_id": task_id,
"deleted_database": {},
"deleted_disk": disk_result,
}
async def _get_task(self, db: AsyncSession, task_id: str) -> Optional[SystemTaskORM]:
result = await db.execute(select(SystemTaskORM).where(SystemTaskORM.task_id == str(task_id or "").strip()))
return result.scalar_one_or_none()
async def _get_jobs(self, db: AsyncSession, task_id: str) -> List[SystemJobORM]:
result = await db.execute(
select(SystemJobORM)
.where(SystemJobORM.task_id == str(task_id or "").strip())
.order_by(SystemJobORM.updated_at.desc(), SystemJobORM.id.desc())
)
return list(result.scalars().all())
async def _get_production_run_for_task(
self,
db: AsyncSession,
task: SystemTaskORM,
jobs: List[SystemJobORM],
) -> Optional[DinsarProductionRunORM]:
candidates: List[str] = []
params = task.params or {}
if isinstance(params, dict):
for key in ("production_run_id", "run_id"):
if params.get(key):
candidates.append(str(params[key]))
for job in jobs:
payload = job.payload or {}
if isinstance(payload, dict):
for key in ("production_run_id", "run_id", "dinsar_run_id"):
if payload.get(key):
candidates.append(str(payload[key]))
if candidates:
result = await db.execute(select(DinsarProductionRunORM).where(DinsarProductionRunORM.run_id.in_(candidates)))
run = result.scalar_one_or_none()
if run is not None:
return run
result = await db.execute(select(DinsarProductionRunORM).where(DinsarProductionRunORM.task_id == task.task_id))
return result.scalar_one_or_none()
async def _build_task_summary(self, db: AsyncSession, task: SystemTaskORM) -> Dict[str, Any]:
jobs = await self._get_jobs(db, task.task_id)
run = await self._get_production_run_for_task(db, task, jobs)
item_counts: Dict[str, int] = {}
execution_counts: Dict[str, int] = {}
if run is not None:
item_counts = await self._status_counts(db, DinsarProductionRunItemORM, run.run_id)
execution_counts = await self._status_counts(db, DinsarProductionExecutionORM, run.run_id)
findings, cleanup_supported, blockers = self._diagnose_findings(
task=task,
jobs=jobs,
run=run,
item_counts=item_counts,
execution_counts=execution_counts,
disk_paths=[],
)
log_count = await self._count(db, select(func.count()).select_from(TaskLogORM).where(TaskLogORM.task_id == task.task_id))
return {
**self._task_payload(task),
"run_id": run.run_id if run else None,
"job_id": jobs[0].job_id if jobs else None,
"job_status": jobs[0].status if jobs else None,
"issue_level": "danger" if _norm_status(task.status) == "FAILED" else ("warning" if findings else "info"),
"issue_summary": findings[0] if findings else _compact(task.message, 160),
"cleanup_supported": cleanup_supported,
"cleanup_blocked_reason": "; ".join(blockers) if blockers else "",
"counts": {
"jobs": len(jobs),
"logs": log_count,
"run_items": sum(item_counts.values()),
"executions": sum(execution_counts.values()),
"completed_items": item_counts.get("COMPLETED", 0),
"failed_items": item_counts.get("FAILED", 0),
"pending_items": item_counts.get("PENDING", 0),
"running_items": item_counts.get("RUNNING", 0),
},
}
async def _status_counts(self, db: AsyncSession, model: Any, run_id: str) -> Dict[str, int]:
result = await db.execute(
select(model.status, func.count())
.where(model.run_id == run_id)
.group_by(model.status)
)
return {str(status or "UNKNOWN").upper(): int(count or 0) for status, count in result.all()}
async def _recent_logs(self, db: AsyncSession, task_id: str) -> List[Dict[str, Any]]:
result = await db.execute(
select(TaskLogORM)
.where(TaskLogORM.task_id == task_id)
.order_by(TaskLogORM.timestamp.desc(), TaskLogORM.id.desc())
.limit(30)
)
return [
{
"id": item.id,
"level": item.log_level,
"message": _compact(item.message, 500),
"timestamp": _dt(item.timestamp),
}
for item in result.scalars().all()
]
async def _collect_run_disk_paths(self, db: AsyncSession, run: DinsarProductionRunORM) -> List[Dict[str, Any]]:
paths: Dict[str, Dict[str, Any]] = {}
result = await db.execute(select(DinsarProductionExecutionORM).where(DinsarProductionExecutionORM.run_id == run.run_id))
for execution in result.scalars().all():
landsar_work_dir = self._landsar_work_dir_for_execution(run, execution)
if landsar_work_dir:
payload = self._path_payload(landsar_work_dir, "landsar_work")
payload["run_key"] = execution.run_key
payload["execution_status"] = execution.status
paths[_normalize_path_text(landsar_work_dir).lower()] = payload
return list(paths.values())
def _landsar_work_dir_for_execution(
self,
run: DinsarProductionRunORM,
execution: DinsarProductionExecutionORM,
) -> str:
if _norm_status(run.engine_code) != "LANDSAR":
return ""
if _norm_status(run.status) in ACTIVE_TASK_STATUSES:
return ""
if _norm_status(execution.status) not in LANDSAR_WORK_TERMINAL_EXECUTION_STATUSES:
return ""
work_root = _normalize_path_text(settings.LANDSAR_WORK_ROOT)
run_key = str(execution.run_key or "").strip()
if not work_root or not run_key or Path(run_key).name != run_key or not run_key.startswith("run_"):
return ""
candidate = _normalize_path_text(os.path.join(work_root, run_key))
return candidate if self._is_landsar_work_delete_path(candidate) else ""
def _diagnose_findings(
self,
*,
task: SystemTaskORM,
jobs: List[SystemJobORM],
run: Optional[DinsarProductionRunORM],
item_counts: Dict[str, int],
execution_counts: Dict[str, int],
disk_paths: List[Dict[str, Any]],
) -> tuple[List[str], bool, List[str]]:
findings: List[str] = []
blockers: List[str] = []
status = _norm_status(task.status)
if status == "FAILED":
findings.append("任务已失败,需要人工确认后清理。")
elif status == "PARTIAL_SUCCESS":
findings.append("任务部分成功,可按需清理 LandSAR 工作目录。")
elif status == "CANCELLED":
findings.append("任务已取消,可按需清理 LandSAR 工作目录。")
elif status in ACTIVE_TASK_STATUSES:
blockers.append("任务仍处于活动状态。")
active_jobs = [job for job in jobs if _norm_status(job.status) in ACTIVE_JOB_STATUSES]
stale_cutoff = _utcnow_naive() - timedelta(seconds=int(getattr(settings, "JOB_WORKER_STALE_RUNNING_SECONDS", 7200)))
stale_jobs = [
job for job in active_jobs
if job.heartbeat_at is not None and job.heartbeat_at < stale_cutoff
]
if stale_jobs:
findings.append(f"发现 {len(stale_jobs)} 个心跳超时 job。")
if active_jobs and not stale_jobs:
blockers.append("仍存在活动 job。")
residual_running = item_counts.get("RUNNING", 0) + execution_counts.get("RUNNING", 0)
if run is not None and residual_running:
findings.append(f"生产 run 中仍有 {residual_running} 个 RUNNING 残留。")
if run is not None and item_counts.get("PENDING", 0):
findings.append(f"生产 run 中仍有 {item_counts.get('PENDING', 0)} 个 PENDING 项未执行。")
missing_paths = [item for item in disk_paths if item.get("path") and not item.get("exists")]
if missing_paths:
findings.append(f"有 {len(missing_paths)} 个登记路径已不存在。")
landsar_work_paths = [
item
for item in disk_paths
if item.get("kind") == "landsar_work" and item.get("exists")
]
if landsar_work_paths:
findings.append(f"发现 {len(landsar_work_paths)} 个可清理的 LandSAR 工作目录。")
cleanup_supported = (
_norm_status(task.task_type) in SUPPORTED_CLEANUP_TASK_TYPES
and not blockers
and bool(landsar_work_paths)
)
return findings, cleanup_supported, blockers
async def _count(self, db: AsyncSession, stmt: Any) -> int:
return int((await db.execute(stmt)).scalar_one() or 0)
def _cleanup_disk_targets(self, diagnosis: Dict[str, Any]) -> List[Dict[str, Any]]:
targets: Dict[str, Dict[str, Any]] = {}
for item in diagnosis.get("disk_paths") or []:
if item.get("kind") != "landsar_work":
continue
path = _normalize_path_text(item.get("path"))
if not path:
continue
payload = self._path_payload(path, item.get("kind") or "unknown")
payload["allowed"] = self._is_allowed_delete_path(path)
targets[path.lower()] = payload
return list(targets.values())
def _path_payload(self, path: str, kind: str) -> Dict[str, Any]:
normalized = _normalize_path_text(path)
exists = _path_exists(normalized)
return {
"path": normalized,
"kind": kind,
"exists": exists,
"allowed": self._is_allowed_delete_path(normalized),
}
def _is_allowed_delete_path(self, path: str) -> bool:
normalized = _normalize_path_text(path)
if not normalized:
return False
return self._is_landsar_work_delete_path(normalized)
def _is_landsar_work_delete_path(self, path: str) -> bool:
work_root = _normalize_path_text(settings.LANDSAR_WORK_ROOT)
normalized = _normalize_path_text(path)
if not work_root or not normalized:
return False
root_full = os.path.abspath(work_root)
full = os.path.abspath(normalized)
if full == root_full or not _is_within_path(full, root_full):
return False
try:
relative = os.path.relpath(full, root_full)
except ValueError:
return False
parts = [part for part in relative.split(os.sep) if part]
return len(parts) == 1 and parts[0].startswith("run_")
async def _delete_disk_targets(self, preview: Dict[str, Any], options: Dict[str, bool]) -> Dict[str, Any]:
delete_landsar_work_dir = bool(options.get("delete_landsar_work_dir", True))
deleted: List[str] = []
missing: List[str] = []
skipped: List[str] = []
failed: List[Dict[str, str]] = []
for item in preview.get("disk_deletes") or []:
kind = item.get("kind")
path = _normalize_path_text(item.get("path"))
if kind != "landsar_work":
skipped.append(path)
continue
if not item.get("allowed"):
skipped.append(path)
continue
if not delete_landsar_work_dir:
skipped.append(path)
continue
if not _path_exists(path):
missing.append(path)
continue
try:
if os.path.isdir(path):
shutil.rmtree(path)
else:
Path(path).unlink()
deleted.append(path)
except OSError as exc:
failed.append({"path": path, "error": str(exc)})
return {
"deleted": deleted,
"missing": missing,
"skipped": skipped,
"failed": failed,
}
def _task_payload(self, task: SystemTaskORM) -> Dict[str, Any]:
return {
"id": task.id,
"task_id": task.task_id,
"task_type": task.task_type,
"task_name": task.task_name,
"status": task.status,
"progress": task.progress,
"message": _compact(task.message),
"params": task.params,
"created_at": _dt(task.created_at),
"updated_at": _dt(task.updated_at),
"started_at": _dt(task.started_at),
"ended_at": _dt(task.ended_at),
}
def _job_payload(self, job: SystemJobORM) -> Dict[str, Any]:
return {
"job_id": job.job_id,
"job_type": job.job_type,
"status": job.status,
"attempts": job.attempts,
"max_attempts": job.max_attempts,
"locked_by": job.locked_by,
"locked_at": _dt(job.locked_at),
"heartbeat_at": _dt(job.heartbeat_at),
"started_at": _dt(job.started_at),
"finished_at": _dt(job.finished_at),
"last_error": _compact(job.last_error),
}
def _run_payload(self, run: DinsarProductionRunORM) -> Dict[str, Any]:
return {
"run_id": run.run_id,
"task_id": run.task_id,
"workflow_run_id": run.workflow_run_id,
"engine_code": run.engine_code,
"profile_code": run.profile_code,
"mode": run.mode,
"source_root": run.source_root,
"publish_root_dir": run.publish_root_dir,
"status": run.status,
"total_items": run.total_items,
"completed_items": run.completed_items,
"failed_items": run.failed_items,
"skipped_items": run.skipped_items,
"latest_message": _compact(run.latest_message),
"created_at": _dt(run.created_at),
"updated_at": _dt(run.updated_at),
"started_at": _dt(run.started_at),
"ended_at": _dt(run.ended_at),
}
ops_maintenance_service = OpsMaintenanceService()