整合求解器活动监控与步长回归证据
同步远端 PNL0003 诊断和大采样网格能力,语义合并活动感知的 60 秒真停滞判定与旧后端 15 分钟兼容兜底。 纳管热路径优化、15 单元运行证据、浏览器与 API 报告,并补充北京时间更新日志和遗留问题。
This commit is contained in:
1 parent
c19cf77aee
commit
e18399c022
46 files changed
+181589
-170
No files matched your search
+55
-21
@@ -3,7 +3,7 @@ from __future__ import annotations
|
||||
from collections.abc import AsyncIterator, Callable, Iterator, Mapping
|
||||
from contextlib import asynccontextmanager
|
||||
import csv
|
||||
from dataclasses import dataclass
|
||||
from dataclasses import dataclass, field
|
||||
from datetime import datetime, timezone
|
||||
import io
|
||||
import json
|
||||
@@ -24,6 +24,7 @@ from pydantic import BaseModel, ConfigDict, Field, ValidationError
|
||||
|
||||
from app.simulation.performance import performance_span, profile_phase, profile_run
|
||||
from app.simulation.property_cache import property_cache_run
|
||||
from app.simulation.solvers.solver import SolverActivityTracker
|
||||
from app.system_xml import (
|
||||
SystemXmlDocument,
|
||||
SystemXmlValidationReport,
|
||||
@@ -86,6 +87,9 @@ SimulationTaskStatus = Literal[
|
||||
class SimulationTaskRecord:
|
||||
simulation_id: str
|
||||
cancel_event: threading.Event
|
||||
activity_tracker: SolverActivityTracker = field(
|
||||
default_factory=SolverActivityTracker
|
||||
)
|
||||
status: SimulationTaskStatus = "queued"
|
||||
cancel_reason: SimulationCancelReason | None = None
|
||||
result: dict[str, object] | None = None
|
||||
@@ -574,6 +578,7 @@ def _simulation_task_snapshot(task: SimulationTaskRecord) -> dict[str, object]:
|
||||
"cancelReason": task.cancel_reason,
|
||||
"result": task.result,
|
||||
"error": task.error,
|
||||
**task.activity_tracker.snapshot().as_dict(),
|
||||
}
|
||||
|
||||
|
||||
@@ -681,6 +686,7 @@ def run_system_xml_simulation(
|
||||
xml_bytes: bytes,
|
||||
progress_callback: SimulationProgressEmitter | None = None,
|
||||
cancel_check: Callable[[], bool] | None = None,
|
||||
activity_tracker: SolverActivityTracker | None = None,
|
||||
) -> dict[str, object]:
|
||||
with property_cache_run() as property_cache:
|
||||
with profile_run() as trace:
|
||||
@@ -688,6 +694,7 @@ def run_system_xml_simulation(
|
||||
xml_bytes,
|
||||
progress_callback,
|
||||
cancel_check,
|
||||
activity_tracker,
|
||||
)
|
||||
|
||||
performance = trace.snapshot()
|
||||
@@ -712,6 +719,7 @@ def _run_system_xml_simulation_profiled(
|
||||
xml_bytes: bytes,
|
||||
progress_callback: SimulationProgressEmitter | None = None,
|
||||
cancel_check: Callable[[], bool] | None = None,
|
||||
activity_tracker: SolverActivityTracker | None = None,
|
||||
) -> dict[str, object]:
|
||||
from app.simulation.solvers.algebraic import AlgebraicSolveError
|
||||
from app.simulation.solvers.solver import SolveIVPConfig
|
||||
@@ -774,6 +782,7 @@ def _run_system_xml_simulation_profiled(
|
||||
sample_step=document.simulation.sample_step,
|
||||
progress_callback=report_system_progress,
|
||||
cancel_check=cancel_check,
|
||||
activity_tracker=activity_tracker,
|
||||
)
|
||||
except SimulationPreparationError as exc:
|
||||
raise HTTPException(
|
||||
@@ -862,6 +871,9 @@ def simulation_event_stream(
|
||||
latest_message = "正在等待仿真任务启动"
|
||||
latest_simulated_time: float | None = None
|
||||
latest_total_time: float | None = None
|
||||
activity_tracker = (
|
||||
task.activity_tracker if task is not None else SolverActivityTracker()
|
||||
)
|
||||
|
||||
def emit_progress(
|
||||
progress: int,
|
||||
@@ -889,6 +901,7 @@ def simulation_event_stream(
|
||||
event["simulatedTime"] = latest_simulated_time
|
||||
if latest_total_time is not None:
|
||||
event["totalTime"] = latest_total_time
|
||||
event.update(activity_tracker.snapshot().as_dict())
|
||||
events.put(event)
|
||||
|
||||
def worker() -> None:
|
||||
@@ -899,28 +912,44 @@ def simulation_event_stream(
|
||||
xml_bytes,
|
||||
emit_progress,
|
||||
task.cancel_event.is_set if task is not None else None,
|
||||
activity_tracker,
|
||||
)
|
||||
if task is not None:
|
||||
result = _mark_simulation_task_result(task, result)
|
||||
result_status = str(result.get("status", "completed"))
|
||||
final_simulated_time = result.get("simulatedUntil")
|
||||
final_activity_kind = (
|
||||
"complete" if result_status == "completed" else result_status
|
||||
)
|
||||
if activity_tracker.snapshot().activity_kind != final_activity_kind:
|
||||
activity_tracker.record_phase(
|
||||
final_activity_kind,
|
||||
(
|
||||
float(final_simulated_time)
|
||||
if isinstance(final_simulated_time, (int, float))
|
||||
and isfinite(final_simulated_time)
|
||||
else None
|
||||
),
|
||||
)
|
||||
result_messages = {
|
||||
"completed": "仿真完成",
|
||||
"stopped": "仿真已由用户终止,已保留部分结果",
|
||||
"stalled": "仿真因进度连接异常而终止,已保留部分结果",
|
||||
"failed": "仿真异常终止,已保留可用的部分结果",
|
||||
}
|
||||
events.put(
|
||||
{
|
||||
"event": "result",
|
||||
"progress": 100 if result_status == "completed" else latest_progress,
|
||||
"phase": result_status,
|
||||
"message": result_messages.get(result_status, "仿真任务结束"),
|
||||
"simulatedTime": result.get("simulatedUntil"),
|
||||
"totalTime": result.get("requestedStopTime"),
|
||||
"result": result,
|
||||
}
|
||||
)
|
||||
result_event = {
|
||||
"event": "result",
|
||||
"progress": 100 if result_status == "completed" else latest_progress,
|
||||
"phase": result_status,
|
||||
"message": result_messages.get(result_status, "仿真任务结束"),
|
||||
"simulatedTime": result.get("simulatedUntil"),
|
||||
"totalTime": result.get("requestedStopTime"),
|
||||
"result": result,
|
||||
}
|
||||
result_event.update(activity_tracker.snapshot().as_dict())
|
||||
events.put(result_event)
|
||||
except HTTPException as exc:
|
||||
activity_tracker.record_phase("failed")
|
||||
detail = exc.detail
|
||||
message = (
|
||||
str(detail.get("message", "仿真失败"))
|
||||
@@ -935,10 +964,12 @@ def simulation_event_stream(
|
||||
"message": message,
|
||||
"detail": detail,
|
||||
}
|
||||
error_event.update(activity_tracker.snapshot().as_dict())
|
||||
if task is not None:
|
||||
_mark_simulation_task_error(task, error_event)
|
||||
events.put(error_event)
|
||||
except Exception as exc: # pragma: no cover - last-resort stream guard
|
||||
activity_tracker.record_phase("failed")
|
||||
error_event = {
|
||||
"event": "error",
|
||||
"progress": latest_progress,
|
||||
@@ -947,6 +978,7 @@ def simulation_event_stream(
|
||||
"message": "仿真服务发生未预期错误。",
|
||||
"detail": str(exc),
|
||||
}
|
||||
error_event.update(activity_tracker.snapshot().as_dict())
|
||||
if task is not None:
|
||||
_mark_simulation_task_error(task, error_event)
|
||||
events.put(error_event)
|
||||
@@ -964,16 +996,18 @@ def simulation_event_stream(
|
||||
try:
|
||||
event = events.get(timeout=SIMULATION_STREAM_HEARTBEAT_SECONDS)
|
||||
except queue.Empty:
|
||||
heartbeat_event = {
|
||||
"event": "progress",
|
||||
"progress": latest_progress,
|
||||
"phase": latest_phase,
|
||||
"message": latest_message,
|
||||
"heartbeat": True,
|
||||
"simulatedTime": latest_simulated_time,
|
||||
"totalTime": latest_total_time,
|
||||
}
|
||||
heartbeat_event.update(activity_tracker.snapshot().as_dict())
|
||||
yield json.dumps(
|
||||
{
|
||||
"event": "progress",
|
||||
"progress": latest_progress,
|
||||
"phase": latest_phase,
|
||||
"message": latest_message,
|
||||
"heartbeat": True,
|
||||
"simulatedTime": latest_simulated_time,
|
||||
"totalTime": latest_total_time,
|
||||
},
|
||||
heartbeat_event,
|
||||
ensure_ascii=False,
|
||||
separators=(",", ":"),
|
||||
) + "\n"
|
||||
|
||||
Reference in new issue
Block a user