140 lines
6.7 KiB
Python
140 lines
6.7 KiB
Python
"""Nonblocking startup diagnostics; editor availability is independent of GCC."""
|
|
from __future__ import annotations
|
|
|
|
from copy import deepcopy
|
|
from datetime import datetime, timezone
|
|
import logging
|
|
import os
|
|
from threading import Event, Lock, Thread
|
|
from time import perf_counter
|
|
from uuid import uuid4
|
|
|
|
from .native_codegen.self_test import CheckCancelled, CheckFailure, check_native_runtime
|
|
|
|
LOGGER = logging.getLogger(__name__)
|
|
MAX_RETRIES = 5
|
|
RETRY_DELAY_SECONDS = 0.5
|
|
|
|
|
|
def simulation_warmup_enabled() -> bool:
|
|
value = os.getenv("SIMULATIONAPP_WARMUP", "on").strip().lower()
|
|
if value in {"", "1", "true", "yes", "on"}:
|
|
return True
|
|
if value in {"0", "false", "no", "off"}:
|
|
return False
|
|
raise ValueError("SIMULATIONAPP_WARMUP must be one of: on, off, true, false, 1, 0.")
|
|
|
|
|
|
class SimulationRuntimeCheck:
|
|
"""One bounded background check per application lifespan, with readable stages."""
|
|
|
|
def __init__(self):
|
|
self._lock = Lock()
|
|
self._stop = Event()
|
|
self._thread: Thread | None = None
|
|
self._history: list[dict] = []
|
|
self._started = perf_counter()
|
|
self._report = {"checkId": uuid4().hex, "status": "pending", "stage": "queued",
|
|
"message": "等待仿真环境启动自检", "durationMs": 0.0,
|
|
"attempt": 0, "maxRetries": MAX_RETRIES, "maxAttempts": MAX_RETRIES + 1,
|
|
"retryDelaySeconds": RETRY_DELAY_SECONDS, "startedAt": None, "finishedAt": None,
|
|
"compiler": None, "compilerVersion": None, "sundialsRoot": None,
|
|
"error": None, "failure": None}
|
|
|
|
def snapshot(self) -> dict:
|
|
with self._lock:
|
|
report = deepcopy(self._report)
|
|
if report["status"] in {"running", "retrying"}:
|
|
report["durationMs"] = (perf_counter() - self._started) * 1000
|
|
elif report["status"] in {"completed", "failed", "disabled", "cancelled"}:
|
|
# Preserve first-failure evidence without surfacing details mid-retry.
|
|
report["attemptHistory"] = deepcopy(self._history)
|
|
return report
|
|
|
|
def _update(self, **values):
|
|
with self._lock:
|
|
if values.get("status") in {"completed", "failed", "disabled", "cancelled"}:
|
|
values["durationMs"] = (perf_counter() - self._started) * 1000
|
|
values["finishedAt"] = datetime.now(timezone.utc).isoformat()
|
|
self._report.update(values)
|
|
|
|
def start(self):
|
|
with self._lock:
|
|
if self._thread is not None:
|
|
return
|
|
self._started = perf_counter()
|
|
self._report.update(status="running", message="正在进行仿真环境首次启动自检",
|
|
startedAt=datetime.now(timezone.utc).isoformat())
|
|
self._thread = Thread(target=self._run, name="simulation-runtime-check", daemon=True)
|
|
self._thread.start()
|
|
|
|
def _run(self):
|
|
for attempt in range(1, MAX_RETRIES + 2):
|
|
if self._stop.is_set():
|
|
self._cancel()
|
|
return
|
|
attempt_started = perf_counter()
|
|
self._update(status="running", attempt=attempt, stage="queued", error=None, failure=None,
|
|
compiler=None, compilerVersion=None, sundialsRoot=None)
|
|
|
|
def progress(**values):
|
|
if attempt > 1 and "message" in values:
|
|
values["message"] = f"正在重试(第 {attempt - 1}/{MAX_RETRIES} 轮):{values['message']}"
|
|
self._update(**values)
|
|
|
|
try:
|
|
if not simulation_warmup_enabled():
|
|
self._update(status="disabled", message="仿真环境启动自检已由服务器配置关闭")
|
|
return
|
|
check_native_runtime(progress, self._stop)
|
|
if self._stop.is_set():
|
|
raise CheckCancelled()
|
|
except CheckCancelled:
|
|
self._cancel()
|
|
return
|
|
except Exception as exc:
|
|
failure = deepcopy(exc.details) if isinstance(exc, CheckFailure) else {
|
|
"command": [], "cwd": None, "exitCode": None, "stdout": "", "stderr": "",
|
|
"errorType": type(exc).__name__,
|
|
}
|
|
with self._lock:
|
|
self._history.append({"attempt": attempt, "status": "failed", "stage": self._report["stage"],
|
|
"durationMs": (perf_counter() - attempt_started) * 1000,
|
|
"error": str(exc), "failure": failure})
|
|
if self._stop.is_set():
|
|
self._cancel()
|
|
return
|
|
if attempt <= MAX_RETRIES:
|
|
previous = "首次检测" if attempt == 1 else f"第 {attempt - 1} 轮重试"
|
|
message = (f"{previous}未通过,正在准备第 {attempt}/{MAX_RETRIES} 轮重试"
|
|
f"(间隔 {RETRY_DELAY_SECONDS:g} 秒);编辑器可正常使用")
|
|
self._update(status="retrying", stage="retry-wait", message=message)
|
|
LOGGER.warning("%s", message)
|
|
# Interruptible on both Windows and Linux; never sleep on the API thread.
|
|
if self._stop.wait(RETRY_DELAY_SECONDS):
|
|
self._cancel()
|
|
return
|
|
continue
|
|
self._update(status="failed", error=str(exc), failure=failure,
|
|
message=f"本次启动自检未通过:首次检测及 {MAX_RETRIES} 轮重试均失败;编辑器仍可使用")
|
|
LOGGER.error("Simulation startup self-check exhausted retries at %s; attempt history: %s",
|
|
self.snapshot()["stage"], self._history, exc_info=True)
|
|
return
|
|
else:
|
|
with self._lock:
|
|
self._history.append({"attempt": attempt, "status": "completed", "stage": "complete",
|
|
"durationMs": (perf_counter() - attempt_started) * 1000,
|
|
"error": None, "failure": None})
|
|
recovery = f"(第 {attempt - 1} 轮重试成功)" if attempt > 1 else ""
|
|
self._update(status="completed", stage="complete", error=None, failure=None,
|
|
message=f"本次启动自检通过,仿真环境可用{recovery}")
|
|
return
|
|
|
|
def _cancel(self):
|
|
self._update(status="cancelled", message="后端正在关闭,启动自检及重试已取消")
|
|
|
|
def close(self):
|
|
self._stop.set()
|
|
if self._thread is not None:
|
|
self._thread.join(timeout=1)
|