优化仿真求解性能并修复流量闭合问题(初版)

This commit is contained in:
ljz committed 2026-08-16 17:46:05 +08:00
1 parent 57b459bc72
commit 5332a788f3
55 files changed
+8973 -549

No files matched your search

+37 -9
View File
@@ -1,6 +1,7 @@
from __future__ import annotations
from collections.abc import Callable, Iterator, Mapping
from collections.abc import AsyncIterator, Callable, Iterator, Mapping
from contextlib import asynccontextmanager
import csv
from dataclasses import dataclass
from datetime import datetime, timezone
@@ -22,6 +23,7 @@ from fastapi.responses import FileResponse, HTMLResponse, StreamingResponse
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.system_xml import (
SystemXmlDocument,
SystemXmlValidationReport,
@@ -35,7 +37,19 @@ if TYPE_CHECKING:
from app.simulation.systems.network import SimulationNetwork
app = FastAPI(title="System Simulation ReactFlow App")
@asynccontextmanager
async def _app_lifespan(application: FastAPI) -> AsyncIterator[None]:
from app.simulation.warmup import warm_up_simulation_runtime
application.state.simulation_warmup = warm_up_simulation_runtime().as_dict()
yield
app = FastAPI(
title="System Simulation ReactFlow App",
lifespan=_app_lifespan,
)
FRONTEND_DIST_DIR = Path(__file__).resolve().parent.parent / "frontend" / "dist"
PROJECT_STORAGE_DIR = Path(__file__).parent / "data" / "reactflow-projects"
SYSTEM_XML_SCHEMA_VERSION = "3"
@@ -668,14 +682,25 @@ def run_system_xml_simulation(
progress_callback: SimulationProgressEmitter | None = None,
cancel_check: Callable[[], bool] | None = None,
) -> dict[str, object]:
with profile_run() as trace:
result = _run_system_xml_simulation_profiled(
xml_bytes,
progress_callback,
cancel_check,
)
with property_cache_run() as property_cache:
with profile_run() as trace:
result = _run_system_xml_simulation_profiled(
xml_bytes,
progress_callback,
cancel_check,
)
performance = trace.snapshot()
if performance.get("mode") == "audit" and property_cache is not None:
cache_info = property_cache.info()
performance["propertyCache"] = {
"hits": cache_info.hits,
"misses": cache_info.misses,
"maxEntriesPerCache": cache_info.max_entries_per_cache,
"cacheCount": cache_info.cache_count,
"currentEntries": cache_info.current_entries,
"evictions": cache_info.evictions,
}
if performance.get("mode") != "off":
diagnostics = dict(result.get("diagnostics", {}))
diagnostics["performance"] = performance
@@ -766,6 +791,9 @@ def _run_system_xml_simulation_profiled(
},
) from exc
except AlgebraicSolveError as exc:
algebraic_diagnostics = exc.diagnostics.as_dict()
algebraic_diagnostics["scopeKind"] = exc.scope_kind
algebraic_diagnostics["scopeComponents"] = list(exc.scope_components)
raise HTTPException(
status_code=422,
detail={
@@ -778,7 +806,7 @@ def _run_system_xml_simulation_profiled(
"message": str(exc),
}
],
"diagnostics": exc.diagnostics.as_dict(),
"diagnostics": algebraic_diagnostics,
},
) from exc
except StreamSolveError as exc:
+14
View File
@@ -66,6 +66,16 @@ RESULT_VARIABLES / DISPLAY / create()`,再把类路径加入库清单。完整
只统计低频的大阶段;`audit` 才展开 RHS、代数闭合、stream 和物性调用,开销也
明显更高。最终优化收益必须在 `off` 下复测。
Peng–Robinson 氦气的高开销物性默认使用一次仿真内独立的精确 LRU 缓存;不同
仿真任务不会共享条目,仿真结束后自动释放。可在启动进程前设置
`SIMULATIONAPP_PROPERTY_CACHE=off` 做数值和性能 A/B,正常运行保持默认 `on`。
缓存只复用完全相同的输入,不做四舍五入或容差匹配。
FastAPI worker 默认在 lifespan 启动阶段预热 SciPy 积分、非线性求解、稀疏
Jacobian 和 System XML XSD,完成后才开始接收请求。它不会运行业务模型,也不
写入文件;如需诊断冷启动,可设置 `SIMULATIONAPP_WARMUP=off`。每个 worker 都会
独立暖机一次。
```powershell
.venv-win\Scripts\python.exe -m app.simulation.benchmark_performance `
--mode audit --warmups 1 --runs 3 `
@@ -73,6 +83,10 @@ RESULT_VARIABLES / DISPLAY / create()`,再把类路径加入库清单。完整
--output app/data/performance-evaluations/helium-step.json
```
缓存关闭对照可在同一命令中增加 `--disable-property-cache`。缓存容量、命中、
未命中和驱逐数会在 audit 响应的
`diagnostics.performance.propertyCache` 中返回。
基准原始 JSON 默认放到已忽略的 `app/data/` 下。指标字段、实测结果和使用边界见
[`仿真性能评估 2026-08-15`](../../docs/仿真性能评估-2026-08-15.md)。
+6 -24
View File
@@ -58,21 +58,6 @@ def _load_factory_xml(specification: str) -> bytes:
return build_reactflow_system_xml(value)
def _clear_property_caches() -> None:
from app.simulation.components.amesim.media.mediums import (
AmesimHeliumPengRobinsonMedium,
)
for method_name in (
"temperature_from_pressure_enthalpy",
"properties_from_mU",
):
method = getattr(AmesimHeliumPengRobinsonMedium, method_name)
cache_clear = getattr(method, "cache_clear", None)
if cache_clear is not None:
cache_clear()
def _serialize_result_event(result: dict[str, object]) -> bytes:
"""Render the final NDJSON payload shape used by the streaming endpoint."""
@@ -98,15 +83,12 @@ def _run_case(
warmups: int,
runs: int,
cancellable_path: bool,
clear_property_cache: bool,
allow_failures: bool,
) -> dict[str, object]:
from app.main import run_system_xml_simulation
cancel_check = (lambda: False) if cancellable_path else None
for _ in range(warmups):
if clear_property_cache:
_clear_property_caches()
result = run_system_xml_simulation(xml_bytes, cancel_check=cancel_check)
if not bool(result.get("success")) and not allow_failures:
raise RuntimeError(f"Warmup for {name!r} failed: {result.get('message')}")
@@ -118,8 +100,6 @@ def _run_case(
profiles: list[dict[str, object]] = []
final_result: dict[str, object] | None = None
for _ in range(runs):
if clear_property_cache:
_clear_property_caches()
wall_start = perf_counter_ns()
cpu_start = process_time_ns()
result = run_system_xml_simulation(xml_bytes, cancel_check=cancel_check)
@@ -190,9 +170,9 @@ def _parse_arguments(argv: list[str] | None = None) -> argparse.Namespace:
help="Do not pass a cancel callback; use the one-shot SciPy path when eligible.",
)
parser.add_argument(
"--cold-property-cache",
"--disable-property-cache",
action="store_true",
help="Clear the two helium property LRU caches before every warmup and measured run.",
help="Disable the run-local exact property cache for an A/B comparison.",
)
parser.add_argument(
"--allow-failures",
@@ -213,6 +193,9 @@ def _parse_arguments(argv: list[str] | None = None) -> argparse.Namespace:
def main(argv: list[str] | None = None) -> int:
arguments = _parse_arguments(argv)
os.environ["SIMULATIONAPP_PROFILE"] = arguments.mode
os.environ["SIMULATIONAPP_PROPERTY_CACHE"] = (
"off" if arguments.disable_property_cache else "on"
)
cases: list[tuple[str, bytes]] = []
for raw_case in arguments.xml:
@@ -226,7 +209,7 @@ def main(argv: list[str] | None = None) -> int:
"generatedAt": datetime.now(UTC).isoformat(),
"profileMode": arguments.mode,
"cancellableSolverPath": not arguments.direct_path,
"coldPropertyCache": bool(arguments.cold_property_cache),
"propertyCacheEnabled": not arguments.disable_property_cache,
"allowFailures": bool(arguments.allow_failures),
"warmups": arguments.warmups,
"runs": arguments.runs,
@@ -242,7 +225,6 @@ def main(argv: list[str] | None = None) -> int:
warmups=arguments.warmups,
runs=arguments.runs,
cancellable_path=not arguments.direct_path,
clear_property_cache=arguments.cold_property_cache,
allow_failures=arguments.allow_failures,
)
for name, xml_bytes in cases
@@ -19,6 +19,7 @@ class AmesimPnpl01(AlgebraicComponent):
MODEL_TYPE = "amesim_pnpl01"
MODEL_VERSION = "0.1.0"
PRESSURE_FLOW_DEPENDS_ON_STREAM = False
PORTS = (PortDefinition.pneumatic("port_1", nominal_role="bidirectional"),)
PARAMETERS = ()
RESULT_VARIABLES = ()
@@ -58,6 +58,7 @@ class AmesimPnor001(AlgebraicComponent):
MODEL_TYPE = "amesim_pnor001"
MODEL_VERSION = "0.3.0"
PRESSURE_FLOW_DEPENDS_ON_STREAM = True
PORTS = (
PortDefinition.pneumatic("port_1", nominal_role="bidirectional"),
PortDefinition.pneumatic("port_2", nominal_role="bidirectional"),
@@ -462,6 +463,7 @@ class AmesimPnvo001FixedOpening(AlgebraicComponent):
MODEL_TYPE = "amesim_pnvo001_fixed"
MODEL_VERSION = "0.2.0"
PRESSURE_FLOW_DEPENDS_ON_STREAM = True
PORTS = (
PortDefinition.pneumatic("port_2", nominal_role="bidirectional"),
PortDefinition.pneumatic("port_3", nominal_role="bidirectional"),
@@ -698,7 +700,7 @@ class AmesimPnvo001FixedOpening(AlgebraicComponent):
def mass_flow(self, p_2: float, p_3: float) -> float:
if (
isclose(p_2, p_3, rel_tol=1.0e-7, abs_tol=1.0e-9)
isclose(p_2, p_3, rel_tol=0.0, abs_tol=1.0e-8)
or self.effective_area == 0.0
):
return 0.0
@@ -885,6 +887,7 @@ class AmesimPnvo001SignalOpening(AmesimPnvo001FixedOpening):
MODEL_TYPE = "amesim_pnvo001"
MODEL_VERSION = "0.2.0"
PRESSURE_FLOW_DEPENDS_ON_STREAM = True
PORTS = (
PortDefinition.signal("res", nominal_role="input"),
PortDefinition.pneumatic("port_2", nominal_role="bidirectional"),
@@ -49,6 +49,7 @@ class AmesimPnl00r(AlgebraicComponent):
MODEL_TYPE = "amesim_pnl00r"
MODEL_VERSION = "0.3.0"
PRESSURE_FLOW_DEPENDS_ON_STREAM = True
PORTS = (
PortDefinition.pneumatic("port_1", nominal_role="bidirectional"),
PortDefinition.pneumatic("port_2", nominal_role="bidirectional"),
@@ -283,7 +284,7 @@ class AmesimPnl00r(AlgebraicComponent):
return 0.5 * (lower + upper)
def mass_flow(self, p_1: float, p_2: float) -> float:
if isclose(p_1, p_2, rel_tol=1.0e-7, abs_tol=1.0e-9):
if isclose(p_1, p_2, rel_tol=0.0, abs_tol=1.0e-8):
return 0.0
pressure_difference = p_1 - p_2
upstream_pressure = max(p_1, p_2, 1.0)
@@ -910,6 +911,7 @@ class AmesimPnl0002(AmesimPnl0001):
MODEL_TYPE = "amesim_pnl0002"
MODEL_VERSION = "0.6.0"
PRESSURE_FLOW_DEPENDS_ON_STREAM = True
PORTS = (
PortDefinition.pneumatic("port_1", nominal_role="bidirectional"),
PortDefinition.pneumatic("port_2", nominal_role="bidirectional"),
@@ -1396,7 +1398,7 @@ class AmesimPnl0003(DynamicComponent):
port_1 = self._properties(self.state_1)
port_2 = self._properties(self.state_2)
pressure_difference = port_1.p - port_2.p
if isclose(port_1.p, port_2.p, rel_tol=1.0e-7, abs_tol=1.0e-9):
if isclose(port_1.p, port_2.p, rel_tol=0.0, abs_tol=1.0e-8):
return 0.0
upstream = port_1 if pressure_difference > 0.0 else port_2
magnitude = self._mass_flow_for_pressure_drop(
@@ -18,6 +18,7 @@ class _AmesimPneumaticNode(AlgebraicComponent):
balance, matching the AMESim dh2 causality.
"""
PRESSURE_FLOW_DEPENDS_ON_STREAM = False
REFERENCE_PORT = "port_2"
def __init__(self, name: str) -> None:
@@ -122,6 +123,7 @@ class AmesimPn3Node2(_AmesimPneumaticNode):
MODEL_TYPE = "amesim_pn3node2"
MODEL_VERSION = "0.3.0"
PRESSURE_FLOW_DEPENDS_ON_STREAM = False
PORTS = (
PortDefinition.pneumatic("port_1", nominal_role="bidirectional"),
PortDefinition.pneumatic("port_2", nominal_role="bidirectional"),
@@ -158,6 +160,7 @@ class AmesimP4Node2(_AmesimPneumaticNode):
MODEL_TYPE = "amesim_p4node2"
MODEL_VERSION = "0.3.0"
PRESSURE_FLOW_DEPENDS_ON_STREAM = False
PORTS = (
PortDefinition.pneumatic("port_1", nominal_role="bidirectional"),
PortDefinition.pneumatic("port_2", nominal_role="bidirectional"),
@@ -28,6 +28,7 @@ class AmesimPnrp17(AlgebraicComponent):
MODEL_TYPE = "amesim_pnrp17"
MODEL_VERSION = "0.1.0"
PRESSURE_FLOW_DEPENDS_ON_STREAM = False
PORTS = (
PortDefinition.pneumatic("port_1", nominal_role="bidirectional"),
PortDefinition.mechanical_translational("port_2"),
@@ -2,7 +2,6 @@ from __future__ import annotations
from collections.abc import Callable
from dataclasses import dataclass
from functools import lru_cache
from typing import ClassVar
from app.simulation.core.errors import RecoverableTrialStateError
@@ -13,6 +12,7 @@ from app.simulation.core.medium import (
)
from app.simulation.core.peng_robinson import HELIUM_PR, PengRobinsonFluid
from app.simulation.performance import profile_property, record_property_iterations
from app.simulation.property_cache import cache_property_calculation
@dataclass(frozen=True)
@@ -71,6 +71,7 @@ class AmesimHeliumPengRobinsonMedium(IdealGasMedium):
return self.cv
@profile_property("density")
@cache_property_calculation("density")
def density(self, p: float, T: float) -> float:
return self.fluid.density(p, T)
@@ -137,6 +138,7 @@ class AmesimHeliumPengRobinsonMedium(IdealGasMedium):
return factor, exponent
@profile_property("isentropic_density_pressure_factor")
@cache_property_calculation("isentropic_density_pressure_factor")
def isentropic_density_pressure_factor(
self,
p: float,
@@ -205,8 +207,8 @@ class AmesimHeliumPengRobinsonMedium(IdealGasMedium):
h / self.R_gas - self.nasa_enthalpy_constant_K
) / self.nasa_cp_over_R
@profile_property("temperature_from_pressure_enthalpy", track_cache=True)
@lru_cache(maxsize=8192)
@profile_property("temperature_from_pressure_enthalpy")
@cache_property_calculation("temperature_from_pressure_enthalpy")
def temperature_from_pressure_enthalpy(self, p: float, h: float) -> float:
temperature = max(self.temperature_from_enthalpy(h), 2.2)
for _iteration in range(16):
@@ -240,8 +242,8 @@ class AmesimHeliumPengRobinsonMedium(IdealGasMedium):
)
return self.temperature_from_internal_energy(U / m)
@profile_property("properties_from_mU", track_cache=True)
@lru_cache(maxsize=8192)
@profile_property("properties_from_mU")
@cache_property_calculation("properties_from_mU")
def properties_from_mU(
self,
m: float,
@@ -16,6 +16,7 @@ class Orifice(AlgebraicComponent):
MODEL_TYPE = "orifice"
MODEL_VERSION = "1.0.0"
PRESSURE_FLOW_DEPENDS_ON_STREAM = False
PORTS = (
PortDefinition.pneumatic("port_a", nominal_role="inlet"),
PortDefinition.pneumatic("port_b", nominal_role="outlet"),
@@ -16,6 +16,7 @@ class ResistivePipe(AlgebraicComponent):
MODEL_TYPE = "pipe"
MODEL_VERSION = "1.0.0"
PRESSURE_FLOW_DEPENDS_ON_STREAM = False
PORTS = (
PortDefinition.pneumatic("port_a", nominal_role="inlet"),
PortDefinition.pneumatic("port_b", nominal_role="outlet"),
@@ -14,6 +14,7 @@ class Tee(AlgebraicComponent):
MODEL_TYPE = "tee"
MODEL_VERSION = "1.0.0"
PRESSURE_FLOW_DEPENDS_ON_STREAM = False
PORTS = (
PortDefinition.pneumatic("port_in", nominal_role="bidirectional"),
PortDefinition.pneumatic("port_out1", nominal_role="bidirectional"),
+7
View File
@@ -21,6 +21,13 @@ if TYPE_CHECKING:
class Component(ABC):
MODEL_TYPE: ClassVar[str | None] = None
MODEL_VERSION: ClassVar[str | None] = None
# ``True`` means that pressure/flow residuals read values written by
# ``update_stream_outflows`` or ``update_flow_temperature_references``.
# ``False`` is an explicit promise that those residuals are independent of
# stream propagation. ``None`` keeps custom components conservative: when
# they override either stream hook, the closure planner retains the legacy
# full-network thermofluid fixed point.
PRESSURE_FLOW_DEPENDS_ON_STREAM: ClassVar[bool | None] = None
PORTS: ClassVar[tuple[PortDefinition, ...]] = ()
PARAMETERS: ClassVar[tuple[ParameterDefinition, ...]] = ()
RESULT_VARIABLES: ClassVar[tuple[ResultVariableDefinition, ...]] = ()
+24
View File
@@ -654,6 +654,29 @@ def record_property_iterations(
)
def record_property_cache(operation: str, *, hit: bool) -> None:
"""Record one run-local property-cache lookup in audit mode."""
trace = _CURRENT_TRACE.get()
if trace is None or trace.mode != "audit":
return
for frame in reversed(_ACTIVE_SPANS.get()):
if (
frame.property_key is not None
and frame.property_operation == operation
):
layer, medium, _unused_operation = frame.property_key.split("|", 2)
trace._record_cache(
frame.property_key,
operation=operation,
layer=layer,
medium=medium,
hits=1 if hit else 0,
misses=0 if hit else 1,
)
return
__all__ = [
"PROFILE_MODE",
"PerformanceTrace",
@@ -661,5 +684,6 @@ __all__ = [
"profile_phase",
"profile_property",
"profile_run",
"record_property_cache",
"record_property_iterations",
]
+234
View File
@@ -0,0 +1,234 @@
"""Run-local, exact-key cache for expensive thermodynamic calculations.
The cache is deliberately bound to one simulation through ``ContextVar``.
That keeps concurrent runs isolated and releases all cached states when the
run finishes. Keys use the original Python values with no rounding or
tolerance-based reuse that could flatten numerical residuals seen by ODE and
nonlinear solvers.
"""
from __future__ import annotations
from collections.abc import Callable, Generator
from contextlib import contextmanager
from contextvars import ContextVar
from dataclasses import dataclass
from functools import lru_cache, wraps
import os
from typing import ParamSpec, TypeVar
from app.simulation.performance import PROFILE_MODE, record_property_cache
_P = ParamSpec("_P")
_R = TypeVar("_R")
DEFAULT_PROPERTY_CACHE_MAX_ENTRIES = 8192
def _read_cache_enabled() -> bool:
raw_value = os.getenv("SIMULATIONAPP_PROPERTY_CACHE", "on").strip().lower()
if raw_value in {"", "1", "true", "yes", "on"}:
return True
if raw_value in {"0", "false", "no", "off"}:
return False
raise ValueError(
"SIMULATIONAPP_PROPERTY_CACHE must be one of: on, off, true, false, 1, 0."
)
PROPERTY_CACHE_ENABLED = _read_cache_enabled()
@dataclass(frozen=True)
class PropertyCacheInfo:
hits: int
misses: int
max_entries_per_cache: int
cache_count: int
current_entries: int
evictions: int
class SimulationPropertyCache:
"""Bounded C-level LRUs owned by one simulation run."""
def __init__(self, max_entries: int = DEFAULT_PROPERTY_CACHE_MAX_ENTRIES) -> None:
if max_entries <= 0:
raise ValueError("Property cache max_entries must be positive.")
self.max_entries_per_cache = int(max_entries)
self._functions: dict[
tuple[str, int, Callable[..., object]],
Callable[..., object],
] = {}
self._failed_misses: dict[
tuple[str, int, Callable[..., object]],
int,
] = {}
self._owners: dict[int, tuple[object, int]] = {}
self._next_owner_token = 0
def owner_token(self, owner: object) -> int:
"""Return a stable identity token and retain its owner for this run."""
identity = id(owner)
existing = self._owners.get(identity)
if existing is not None and existing[0] is owner:
return existing[1]
self._next_owner_token += 1
self._owners[identity] = (owner, self._next_owner_token)
return self._next_owner_token
def get_or_compute(
self,
operation: str,
owner: object,
function: Callable[..., _R],
args: tuple[object, ...],
kwargs: dict[str, object],
) -> _R:
cache_key = (operation, self.owner_token(owner), function)
cached_function = self._functions.get(cache_key)
if cached_function is None:
@lru_cache(maxsize=self.max_entries_per_cache, typed=True)
def invoke(*cached_args: object, **cached_kwargs: object) -> _R:
return function(owner, *cached_args, **cached_kwargs)
cached_function = invoke
self._functions[cache_key] = cached_function
if PROFILE_MODE != "audit":
try:
return cached_function(*args, **kwargs)
except Exception:
self._failed_misses[cache_key] = (
self._failed_misses.get(cache_key, 0) + 1
)
raise
before = cached_function.cache_info() # type: ignore[attr-defined]
try:
value = cached_function(*args, **kwargs)
except Exception:
self._failed_misses[cache_key] = (
self._failed_misses.get(cache_key, 0) + 1
)
raise
finally:
after = cached_function.cache_info() # type: ignore[attr-defined]
hit = after.hits > before.hits
record_property_cache(operation, hit=hit)
return value
def info(self) -> PropertyCacheInfo:
cache_infos = {
key: cached.cache_info() # type: ignore[attr-defined]
for key, cached in self._functions.items()
}
return PropertyCacheInfo(
hits=sum(info.hits for info in cache_infos.values()),
misses=sum(info.misses for info in cache_infos.values()),
max_entries_per_cache=self.max_entries_per_cache,
cache_count=len(self._functions),
current_entries=sum(info.currsize for info in cache_infos.values()),
evictions=sum(
max(
0,
info.misses
- self._failed_misses.get(key, 0)
- info.currsize,
)
for key, info in cache_infos.items()
),
)
_CURRENT_PROPERTY_CACHE: ContextVar[SimulationPropertyCache | None] = ContextVar(
"simulation_property_cache",
default=None,
)
def current_property_cache() -> SimulationPropertyCache | None:
return _CURRENT_PROPERTY_CACHE.get()
@contextmanager
def property_cache_run(
*,
max_entries: int = DEFAULT_PROPERTY_CACHE_MAX_ENTRIES,
) -> Generator[SimulationPropertyCache | None, None, None]:
"""Bind a fresh cache to one top-level simulation run.
Nested uses reuse the existing cache so lower-level simulation helpers can
safely opt in without replacing the cache created by the API entry point.
"""
existing = _CURRENT_PROPERTY_CACHE.get()
if existing is not None:
yield existing
return
if not PROPERTY_CACHE_ENABLED:
yield None
return
cache = SimulationPropertyCache(max_entries=max_entries)
token = _CURRENT_PROPERTY_CACHE.set(cache)
try:
yield cache
finally:
_CURRENT_PROPERTY_CACHE.reset(token)
def cache_property_calculation(
operation: str,
) -> Callable[[Callable[_P, _R]], Callable[_P, _R]]:
"""Cache one pure property calculation with hashable arguments per run."""
def decorate(function: Callable[_P, _R]) -> Callable[_P, _R]:
if not PROPERTY_CACHE_ENABLED:
return function
@wraps(function)
def wrapper(*args: _P.args, **kwargs: _P.kwargs) -> _R:
cache = _CURRENT_PROPERTY_CACHE.get()
if cache is None:
return function(*args, **kwargs)
owner = args[0] if args else function
return cache.get_or_compute(
operation,
owner,
function,
tuple(args[1:] if args else ()),
dict(kwargs),
)
return wrapper
return decorate
def with_property_cache(function: Callable[_P, _R]) -> Callable[_P, _R]:
"""Ensure a simulation entry point has a run-local cache."""
if not PROPERTY_CACHE_ENABLED:
return function
@wraps(function)
def wrapper(*args: _P.args, **kwargs: _P.kwargs) -> _R:
with property_cache_run():
return function(*args, **kwargs)
return wrapper
__all__ = [
"DEFAULT_PROPERTY_CACHE_MAX_ENTRIES",
"PROPERTY_CACHE_ENABLED",
"PropertyCacheInfo",
"SimulationPropertyCache",
"cache_property_calculation",
"current_property_cache",
"property_cache_run",
"with_property_cache",
]
File diff suppressed because it is too large. Load diff
File diff suppressed because it is too large. Load diff
+45 -19
View File
@@ -3,6 +3,8 @@ from __future__ import annotations
from dataclasses import dataclass
from math import isfinite
from app.simulation.core.base import Component
from app.simulation.core.ports import PortState
from app.simulation.performance import profile_phase
from app.simulation.systems.network import Endpoint, SimulationNetwork
@@ -19,35 +21,62 @@ class PneumaticVolumeDiagnostics:
}
@dataclass(frozen=True)
class _PneumaticVolumeConnectionBinding:
connected_endpoint: Endpoint
connected_port: PortState
class PneumaticVolumeResolver:
"""Propagate AMESim pneumatic external-volume connector variables."""
def __init__(self, network: SimulationNetwork) -> None:
self.network = network
self._pneumatic_ports = tuple(
component.get_port(definition.name)
for component in network.components.values()
for definition in component.active_port_definitions
if definition.kind == "physical" and definition.domain == "pneumatic"
)
self._output_components = tuple(
component
for component in network.components.values()
if type(component).pneumatic_volume_outputs
is not Component.pneumatic_volume_outputs
)
self._connected_endpoint = self._build_connection_map()
self.last_diagnostics: PneumaticVolumeDiagnostics | None = None
def _build_connection_map(self) -> dict[Endpoint, Endpoint]:
result: dict[Endpoint, Endpoint] = {}
def _build_connection_map(
self,
) -> dict[Endpoint, _PneumaticVolumeConnectionBinding]:
result: dict[Endpoint, _PneumaticVolumeConnectionBinding] = {}
for connection in self.network.connections:
if connection.kind != "physical" or connection.domain != "pneumatic":
continue
first, second = connection.endpoints
result[first] = second
result[second] = first
result[first] = _PneumaticVolumeConnectionBinding(
connected_endpoint=second,
connected_port=self.network.components[second.component].get_port(
second.port
),
)
result[second] = _PneumaticVolumeConnectionBinding(
connected_endpoint=first,
connected_port=self.network.components[first.component].get_port(
first.port
),
)
return result
@profile_phase("simulation.pneumatic_volume", minimum_mode="audit")
def solve(self) -> PneumaticVolumeDiagnostics:
for component in self.network.components.values():
for definition in component.active_port_definitions:
if definition.kind == "physical" and definition.domain == "pneumatic":
port = component.get_port(definition.name)
port.volume = 0.0
port.volume_flow = 0.0
for port in self._pneumatic_ports:
port.volume = 0.0
port.volume_flow = 0.0
outputs: dict[Endpoint, tuple[float, float]] = {}
for component in self.network.components.values():
for component in self._output_components:
for port_name, raw_values in component.pneumatic_volume_outputs().items():
port = component.get_port(port_name)
definition = port.definition
@@ -73,18 +102,15 @@ class PneumaticVolumeResolver:
propagated = 0
for endpoint, values in outputs.items():
connected = self._connected_endpoint.get(endpoint)
if connected is None:
binding = self._connected_endpoint.get(endpoint)
if binding is None:
continue
if connected in outputs:
if binding.connected_endpoint in outputs:
raise ValueError(
"A pneumatic connection cannot contain two external-volume "
f"sources: {endpoint} and {connected}."
f"sources: {endpoint} and {binding.connected_endpoint}."
)
connected_port = self.network.components[connected.component].get_port(
connected.port
)
connected_port.volume, connected_port.volume_flow = values
binding.connected_port.volume, binding.connected_port.volume_flow = values
propagated += 1
diagnostics = PneumaticVolumeDiagnostics(
+54 -20
View File
@@ -2,8 +2,10 @@ from __future__ import annotations
from dataclasses import dataclass
from math import isfinite
from typing import Protocol
from typing import Callable, Protocol
from app.simulation.core.base import Component
from app.simulation.core.ports import PortState
from app.simulation.performance import profile_phase
from app.simulation.systems.network import Endpoint, SimulationNetwork
@@ -38,31 +40,56 @@ class SignalSolveDiagnostics:
return {"propagated": self.propagated}
@dataclass(frozen=True)
class _SignalOutputBinding:
component: Component
evaluate: Callable[[float], dict[str, float]]
@dataclass(frozen=True)
class _SignalConnectionBinding:
source: PortState
target: PortState
class SignalResolver:
"""Propagate scalar signal connections from output ports to input ports."""
def __init__(self, network: SimulationNetwork) -> None:
self.network = network
self._connections = [
connection for connection in network.connections if connection.kind == "signal"
]
self._output_bindings = tuple(
_SignalOutputBinding(component=component, evaluate=evaluate)
for component in network.components.values()
if (evaluate := getattr(component, "signal_output_values", None)) is not None
)
self._event_sources = tuple(
(component.name, source_event_times)
for component in network.components.values()
if (
source_event_times := getattr(
component,
"signal_event_times",
None,
)
)
is not None
)
self._connections = tuple(
self._connection_binding(connection.endpoints)
for connection in network.connections
if connection.kind == "signal"
)
self.last_diagnostics: SignalSolveDiagnostics | None = None
@profile_phase("simulation.signal", minimum_mode="audit")
def solve(self, time: float) -> SignalSolveDiagnostics:
for component in self.network.components.values():
signal_output_values = getattr(component, "signal_output_values", None)
if signal_output_values is None:
continue
for port_name, value in signal_output_values(time).items():
component.get_port(port_name).signal = float(value)
for binding in self._output_bindings:
for port_name, value in binding.evaluate(time).items():
binding.component.get_port(port_name).signal = float(value)
propagated = 0
for connection in self._connections:
source, target = self._source_target(connection.endpoints)
source_port = self.network.components[source.component].get_port(source.port)
target_port = self.network.components[target.component].get_port(target.port)
target_port.signal = source_port.signal
for binding in self._connections:
binding.target.signal = binding.source.signal
propagated += 1
diagnostics = SignalSolveDiagnostics(propagated=propagated)
@@ -86,15 +113,12 @@ class SignalResolver:
return ()
events: set[float] = set()
for component in self.network.components.values():
source_event_times = getattr(component, "signal_event_times", None)
if source_event_times is None:
continue
for component_name, source_event_times in self._event_sources:
for raw_time in source_event_times(start, stop):
event_time = float(raw_time)
if not isfinite(event_time):
raise ValueError(
f"Signal event time from component '{component.name}' must be finite."
f"Signal event time from component '{component_name}' must be finite."
)
if start < event_time < stop:
events.add(event_time)
@@ -109,3 +133,13 @@ class SignalResolver:
if second_port.definition is not None and second_port.definition.nominal_role == "output":
return second, first
raise ValueError("Signal connection must contain one output endpoint.")
def _connection_binding(
self,
endpoints: tuple[Endpoint, Endpoint],
) -> _SignalConnectionBinding:
source, target = self._source_target(endpoints)
return _SignalConnectionBinding(
source=self.network.components[source.component].get_port(source.port),
target=self.network.components[target.component].get_port(target.port),
)
+6 -2
View File
@@ -807,9 +807,13 @@ def _integrate_scipy_stepwise(
segment_accepted_steps += 1
step_end_time = float(solver.t)
step_end_state = [float(value) for value in solver.y]
crosses_sample = (
sample_index < len(sample_times)
and sample_times[sample_index] <= step_end_time
)
dense_output = (
solver.dense_output()
if sample_times or state_transition_handler is not None
if crosses_sample or state_transition_handler is not None
else None
)
@@ -918,11 +922,11 @@ def _integrate_scipy_stepwise(
else last_accepted_time
)
if sample_times:
assert dense_output is not None
while (
sample_index < len(sample_times)
and sample_times[sample_index] <= last_accepted_time
):
assert dense_output is not None
sample_time = float(sample_times[sample_index])
sample_state = [
float(value) for value in dense_output(sample_time)
+74 -45
View File
@@ -2,9 +2,10 @@ from __future__ import annotations
from dataclasses import dataclass
from app.simulation.core.base import DynamicComponent
from app.simulation.core.base import Component, DynamicComponent
from app.simulation.core.ports import PortState
from app.simulation.performance import profile_phase
from app.simulation.systems.network import Endpoint, SimulationNetwork
from app.simulation.systems.network import SimulationNetwork
class StreamSolveError(RuntimeError):
@@ -27,6 +28,14 @@ class StreamSolveDiagnostics:
}
@dataclass(frozen=True)
class _StreamConnectionBinding:
component_name: str
port_name: str
connected_component: Component
connected_port: PortState
class StreamResolver:
"""Resolve outflow enthalpy propagation after pressure and flow are known."""
@@ -40,28 +49,59 @@ class StreamResolver:
self.network = network
self.relative_tolerance = relative_tolerance
self.max_iterations = max_iterations
self._connected_endpoint = self._build_connection_map()
self._components = tuple(network.components.values())
self._dynamic_components = tuple(
component
for component in self._components
if isinstance(component, DynamicComponent)
)
self._non_dynamic_components = tuple(
component
for component in self._components
if not isinstance(component, DynamicComponent)
)
self._ports = tuple(
(component.name, port_name, port)
for component in self._components
for port_name, port in component.ports.items()
)
self._connection_bindings = self._build_connection_bindings()
self.last_diagnostics: StreamSolveDiagnostics | None = None
def _build_connection_map(self) -> dict[Endpoint, Endpoint]:
result: dict[Endpoint, Endpoint] = {}
def _build_connection_bindings(self) -> tuple[_StreamConnectionBinding, ...]:
result: list[_StreamConnectionBinding] = []
for connection in self.network.connections:
if connection.kind != "physical":
continue
first, second = connection.endpoints
result[first] = second
result[second] = first
return result
first_component = self.network.components[first.component]
second_component = self.network.components[second.component]
result.append(
_StreamConnectionBinding(
component_name=first.component,
port_name=first.port,
connected_component=second_component,
connected_port=second_component.get_port(second.port),
)
)
result.append(
_StreamConnectionBinding(
component_name=second.component,
port_name=second.port,
connected_component=first_component,
connected_port=first_component.get_port(first.port),
)
)
return tuple(result)
def connected_enthalpies(self) -> dict[str, dict[str, float]]:
values: dict[str, dict[str, float]] = {
component.name: {} for component in self.network.components.values()
component.name: {} for component in self._components
}
for endpoint, connected in self._connected_endpoint.items():
connected_port = self.network.components[connected.component].get_port(
connected.port
for binding in self._connection_bindings:
values[binding.component_name][binding.port_name] = (
binding.connected_port.h_outflow
)
values[endpoint.component][endpoint.port] = connected_port.h_outflow
return values
def connected_temperature_reference_enthalpies(
@@ -70,26 +110,21 @@ class StreamResolver:
"""Return connector references used for upstream temperature only."""
values: dict[str, dict[str, float]] = {
component.name: {} for component in self.network.components.values()
component.name: {} for component in self._components
}
for endpoint, connected in self._connected_endpoint.items():
connected_component = self.network.components[connected.component]
connected_port = connected_component.get_port(connected.port)
values[endpoint.component][endpoint.port] = float(
for binding in self._connection_bindings:
values[binding.component_name][binding.port_name] = float(
getattr(
connected_component,
binding.connected_component,
"temperature_reference_h",
connected_port.h_outflow,
binding.connected_port.h_outflow,
)
)
return values
@profile_phase("simulation.refresh", minimum_mode="audit")
def _refresh_dynamic_components(
self,
components: list[DynamicComponent],
) -> None:
for component in components:
def _refresh_dynamic_components(self) -> None:
for component in self._dynamic_components:
component.refresh_thermodynamic_ports()
@profile_phase("simulation.refresh", minimum_mode="audit")
@@ -97,40 +132,34 @@ class StreamResolver:
self,
connected: dict[str, dict[str, float]],
) -> None:
for component in self.network.components.values():
if isinstance(component, DynamicComponent):
component.refresh_thermodynamic_ports()
else:
component.update_stream_outflows(connected[component.name])
for component in self._non_dynamic_components:
component.update_stream_outflows(connected[component.name])
@profile_phase("simulation.stream", minimum_mode="audit")
def solve(self) -> tuple[StreamSolveDiagnostics, dict[str, dict[str, float]]]:
dynamic_components = [
component
for component in self.network.components.values()
if isinstance(component, DynamicComponent)
]
self._refresh_dynamic_components(dynamic_components)
def solve(
self,
*,
dynamic_ports_are_current: bool = False,
) -> tuple[StreamSolveDiagnostics, dict[str, dict[str, float]]]:
if not dynamic_ports_are_current:
self._refresh_dynamic_components()
max_delta = 0.0
for iteration in range(1, self.max_iterations + 1):
previous = {
(component.name, port_name): port.h_outflow
for component in self.network.components.values()
for port_name, port in component.ports.items()
(component_name, port_name): port.h_outflow
for component_name, port_name, port in self._ports
}
connected = self.connected_enthalpies()
self._refresh_stream_components(connected)
deltas = [
abs(port.h_outflow - previous[(component.name, port_name)])
for component in self.network.components.values()
for port_name, port in component.ports.items()
abs(port.h_outflow - previous[(component_name, port_name)])
for component_name, port_name, port in self._ports
]
magnitudes = [
abs(port.h_outflow)
for component in self.network.components.values()
for port in component.ports.values()
for _component_name, _port_name, port in self._ports
]
max_delta = max(deltas, default=0.0)
scale = max(magnitudes + [1.0])
+576 -22
View File
@@ -5,10 +5,13 @@ from dataclasses import dataclass, replace
from math import floor, isfinite
from typing import Literal
from app.simulation.core.base import DynamicComponent
from app.simulation.core.base import Component, DynamicComponent
from app.simulation.core.metadata import ResultVariableMetadata
from app.simulation.core.ports import PortState
from app.simulation.performance import performance_span, profile_phase
from app.simulation.property_cache import with_property_cache
from app.simulation.solvers.algebraic import PressureFlowSolver
from app.simulation.solvers.algebraic_blocks import StreamPressureBlockSolver
from app.simulation.solvers.mechanical import (
MechanicalConstraintGroup,
MechanicalStateReducer,
@@ -29,6 +32,25 @@ SimulationCancellationCheck = Callable[[], bool]
SimulationRunStatus = Literal["completed", "cancelled", "failed"]
@dataclass(frozen=True)
class _ThermofluidClosurePlan:
"""Static execution data for one compiled network.
The first pressure-flow solve remains global. Later fixed-point passes only
need the physical islands whose constitutive equations read stream-derived
enthalpy. An unclassified custom stream component deliberately falls back
to the original global solve.
"""
physical_ports: tuple[PortState, ...]
global_component_group: tuple[str, ...]
secondary_pressure_solvers: tuple[PressureFlowSolver, ...]
secondary_component_groups: tuple[tuple[str, ...], ...]
uses_conservative_global_solver: bool
conservative_fallback_reason: str | None
secondary_block_solvers: tuple[StreamPressureBlockSolver, ...] = ()
@dataclass(frozen=True)
class SimulationPreparationIssue:
code: str
@@ -357,15 +379,271 @@ class GenericFluidSystem:
self.pneumatic_volume_resolver = PneumaticVolumeResolver(network)
self.signal_resolver = SignalResolver(network)
self.stream_resolver = StreamResolver(network)
self._thermofluid_closure_plan = self._build_thermofluid_closure_plan()
self.algebraic_solve_count = 0
self.algebraic_seeded_solve_count = 0
self.algebraic_nonlinear_solve_count = 0
self.algebraic_optimizer_evaluation_count = 0
self.algebraic_residual_evaluation_count = 0
self.algebraic_block_fallback_count = 0
self.algebraic_dense_fallback_count = 0
self.thermofluid_pressure_pass_count = 0
self.max_algebraic_residual = 0.0
self.max_algebraic_evaluations = 0
self.max_algebraic_residual_evaluations = 0
self._last_algebraic_diagnostics = None
self._last_algebraic_scope: tuple[str, ...] = ()
self.max_stream_iterations = 0
self.max_thermofluid_iterations = 0
self.signal_propagation_count = 0
self.pneumatic_volume_propagation_count = 0
self._jacobian_sparsity = None
def _request_causal_residual_audit(self) -> None:
"""Make topology or mode boundaries verify the next causal closure."""
self.pressure_flow_solver.request_causal_audit()
for block_solver in (
self._thermofluid_closure_plan.secondary_block_solvers
):
block_solver.request_causal_audit()
@staticmethod
def _overrides_stream_update(component: Component) -> bool:
component_type = type(component)
return (
component_type.update_stream_outflows
is not Component.update_stream_outflows
or component_type.update_flow_temperature_references
is not Component.update_flow_temperature_references
)
def _physical_component_groups(self) -> tuple[tuple[str, ...], ...]:
"""Return physical islands in component insertion order."""
physical_names = tuple(
component.name
for component in self.network.components.values()
if any(
definition.kind == "physical"
for definition in component.active_port_definitions
)
)
adjacency = {name: set() for name in physical_names}
for connection in self.network.connections:
if connection.kind != "physical":
continue
first, second = connection.endpoints
adjacency[first.component].add(second.component)
adjacency[second.component].add(first.component)
groups: list[tuple[str, ...]] = []
visited: set[str] = set()
for root in physical_names:
if root in visited:
continue
members = {root}
pending = [root]
visited.add(root)
while pending:
current = pending.pop()
for neighbor in adjacency[current]:
if neighbor in visited:
continue
visited.add(neighbor)
members.add(neighbor)
pending.append(neighbor)
groups.append(tuple(name for name in physical_names if name in members))
return tuple(groups)
def _network_for_component_group(
self,
component_names: tuple[str, ...],
all_physical_names: frozenset[str],
) -> SimulationNetwork:
if frozenset(component_names) == all_physical_names:
return self.network
selected = frozenset(component_names)
subnetwork = SimulationNetwork(
name=f"{self.network.name}:thermofluid:{len(component_names)}"
)
for component in self.network.components.values():
if component.name in selected:
subnetwork.add_component(component)
subnetwork.connections.extend(
connection
for connection in self.network.connections
if connection.kind == "physical"
and connection.endpoint_a.component in selected
and connection.endpoint_b.component in selected
)
return subnetwork
def _pressure_solver_for_component_group(
self,
component_names: tuple[str, ...],
all_physical_names: frozenset[str],
) -> PressureFlowSolver:
subnetwork = self._network_for_component_group(
component_names,
all_physical_names,
)
if subnetwork is self.network:
return self.pressure_flow_solver
return PressureFlowSolver(
subnetwork,
residual_tolerance=self.pressure_flow_solver.residual_tolerance,
max_evaluations=self.pressure_flow_solver.max_evaluations,
scope_kind="physicalIsland",
)
def _build_thermofluid_closure_plan(self) -> _ThermofluidClosurePlan:
physical_ports = tuple(
component.get_port(definition.name)
for component in self.network.components.values()
for definition in component.active_port_definitions
if definition.kind == "physical"
)
physical_groups = self._physical_component_groups()
all_physical_names = frozenset(
name for group in physical_groups for name in group
)
all_physical_order = tuple(
name for group in physical_groups for name in group
)
sensitive_names: set[str] = set()
has_unclassified_stream_component = False
has_invalid_dependency_declaration = False
for name in all_physical_names:
component = self.network.components[name]
# Only an exact-class declaration opts into pruning. A custom
# subclass cannot accidentally inherit a purity promise after
# changing its stream hook or constitutive equations.
declared = type(component).__dict__.get(
"PRESSURE_FLOW_DEPENDS_ON_STREAM"
)
if declared is True:
sensitive_names.add(name)
elif declared is False:
continue
elif declared is None and self._overrides_stream_update(component):
# Preserve the exact legacy behavior for custom components that
# receive stream values but have not declared equation purity.
has_unclassified_stream_component = True
elif declared is not None:
has_invalid_dependency_declaration = True
if has_invalid_dependency_declaration:
return _ThermofluidClosurePlan(
physical_ports=physical_ports,
global_component_group=all_physical_order,
secondary_pressure_solvers=(self.pressure_flow_solver,),
secondary_component_groups=(all_physical_order,),
uses_conservative_global_solver=True,
conservative_fallback_reason="invalidDependencyDeclaration",
)
if has_unclassified_stream_component:
return _ThermofluidClosurePlan(
physical_ports=physical_ports,
global_component_group=all_physical_order,
secondary_pressure_solvers=(self.pressure_flow_solver,),
secondary_component_groups=(all_physical_order,),
uses_conservative_global_solver=True,
conservative_fallback_reason="unclassifiedStreamComponent",
)
component_names = set(self.network.components)
compiled_equations = self.pressure_flow_solver.equation_templates
for equation in compiled_equations:
if equation.owner != "component":
continue
referenced_components = {
parts[0]
for variable in equation.variables
if len(parts := variable.rsplit(".", 2)) == 3
and parts[0] in component_names
}
if referenced_components - {equation.owner_id}:
# Catalog equations are component-local and connectors carry
# cross-component constraints. A custom residual may violate
# that convention, so retain the unsplit global problem.
return _ThermofluidClosurePlan(
physical_ports=physical_ports,
global_component_group=all_physical_order,
secondary_pressure_solvers=(self.pressure_flow_solver,),
secondary_component_groups=(all_physical_order,),
uses_conservative_global_solver=True,
conservative_fallback_reason="crossComponentEquation",
)
for group in physical_groups:
selected = frozenset(group)
connection_ids = {
connection.id
for connection in self.network.connections
if connection.kind == "physical"
and connection.endpoint_a.component in selected
and connection.endpoint_b.component in selected
}
unknown_count = sum(
unknown.component in selected
for unknown in self.pressure_flow_solver.unknowns
)
equation_count = sum(
(
equation.owner == "component"
and equation.owner_id in selected
)
or (
equation.owner == "connection"
and equation.owner_id in connection_ids
)
for equation in compiled_equations
)
if unknown_count != equation_count:
# The full network can be square even when two disconnected
# rectangular islands happen to cancel each other's equation
# count. Preserve the original global least-squares problem in
# that unusual case rather than changing its solution space.
return _ThermofluidClosurePlan(
physical_ports=physical_ports,
global_component_group=all_physical_order,
secondary_pressure_solvers=(self.pressure_flow_solver,),
secondary_component_groups=(all_physical_order,),
uses_conservative_global_solver=True,
conservative_fallback_reason="nonSquarePhysicalIsland",
)
coupled_groups = tuple(
group for group in physical_groups if sensitive_names.intersection(group)
)
secondary_pressure_solvers = tuple(
self._pressure_solver_for_component_group(group, all_physical_names)
for group in coupled_groups
)
secondary_block_solvers = tuple(
StreamPressureBlockSolver(
pressure_solver,
tuple(name for name in group if name in sensitive_names),
)
for pressure_solver, group in zip(
secondary_pressure_solvers,
coupled_groups,
)
)
return _ThermofluidClosurePlan(
physical_ports=physical_ports,
global_component_group=all_physical_order,
secondary_pressure_solvers=secondary_pressure_solvers,
secondary_component_groups=coupled_groups,
uses_conservative_global_solver=False,
conservative_fallback_reason=None,
secondary_block_solvers=secondary_block_solvers,
)
def initial_state_vector(self) -> list[float]:
return self.pneumatic_storage_reducer.synchronize_state_vector(
self.mechanical_state_reducer.initial_state_vector(),
@@ -377,6 +655,129 @@ class GenericFluidSystem:
self.pneumatic_storage_reducer.synchronize_state_vector(values)
)
@staticmethod
def _entry_has_pneumatic_state(entry: object) -> bool:
"""Return whether one reduced ODE entry owns pneumatic state.
Mechanical constraint groups are synthetic state owners. Every other
entry is a dynamic component, so its active port metadata is the
topology-level way to classify it without depending on model names.
"""
if isinstance(entry, MechanicalConstraintGroup):
return False
return any(
definition.kind == "physical" and definition.domain == "pneumatic"
for definition in entry.active_port_definitions
)
def _add_pneumatic_volume_state_dependencies(
self,
dependencies: list[set[int]],
entries: tuple[object, ...],
owner_by_component: dict[str, int],
) -> None:
"""Close the cross-domain dependency hidden by external volume ports.
A pneumatic-volume source such as a piston writes swept volume from
mechanical coordinates into a connected storage component before the
pressure-flow closure. The ordinary physical-path walk intentionally
stops at a storage state. Consequently, a second storage connected to
that chamber can depend on the piston even though the path crosses the
chamber state, and that derivative was previously omitted from the BDF
sparsity pattern.
Reuse the resolver's compiled output/connection plan to locate each
receiving storage. Mechanical states already found from that receiver,
the volume source's own ODE state (when it has one), and pneumatic states
whose local closure reaches the receiver form one conservative
cross-domain dependency set. Add it in both directions. If executable
custom/source metadata cannot bound those drivers, use a dense pattern.
"""
resolver = self.pneumatic_volume_resolver
pneumatic_entries = tuple(
self._entry_has_pneumatic_state(entry) for entry in entries
)
all_entry_indexes = set(range(len(entries)))
def use_conservative_dense_pattern() -> None:
for entry_dependencies in dependencies:
entry_dependencies.update(all_entry_indexes)
for component in resolver._output_components:
# ``pneumatic_volume_outputs`` is executable code rather than an
# equation-level dependency declaration. Catalog components with
# no directed signal input can be bounded by their own ODE state
# and the mechanical states already connected through topology.
# Custom/output components with an external signal driver keep the
# implicit integrator safe by disabling sparsity for this system.
if (
not type(component).__module__.startswith(
"app.simulation.components."
)
or any(
(
definition.kind == "signal"
and definition.nominal_role == "input"
)
or (
definition.kind == "physical"
and definition.domain
not in {"pneumatic", "mechanical"}
)
for definition in component.active_port_definitions
)
):
use_conservative_dense_pattern()
return
source_index = owner_by_component.get(component.name)
receiver_indexes: set[int] = set()
for definition in component.active_port_definitions:
if (
definition.kind != "physical"
or definition.domain != "pneumatic"
):
continue
binding = resolver._connected_endpoint.get(
Endpoint(component.name, definition.name)
)
if binding is None:
continue
receiver_index = owner_by_component.get(
binding.connected_endpoint.component
)
if receiver_index is not None and pneumatic_entries[receiver_index]:
receiver_indexes.add(receiver_index)
for receiver_index in receiver_indexes:
driver_indexes = {
entry_index
for entry_index in dependencies[receiver_index]
if isinstance(entries[entry_index], MechanicalConstraintGroup)
}
if source_index is not None:
driver_indexes.add(source_index)
if not driver_indexes:
use_conservative_dense_pattern()
return
coupled_pneumatic_indexes = {
entry_index
for entry_index, is_pneumatic in enumerate(pneumatic_entries)
if is_pneumatic
and (
entry_index == receiver_index
or receiver_index in dependencies[entry_index]
)
}
for pneumatic_index in coupled_pneumatic_indexes:
dependencies[pneumatic_index].update(driver_indexes)
for driver_index in driver_indexes:
dependencies[driver_index].update(
coupled_pneumatic_indexes
)
def _build_jacobian_sparsity(self):
"""Build a conservative state dependency graph for implicit solvers.
@@ -431,6 +832,12 @@ class GenericFluidSystem:
pending.append(neighbour)
dependencies.append(found)
self._add_pneumatic_volume_state_dependencies(
dependencies,
entries,
owner_by_component,
)
offsets = [0]
for state_size in entry_sizes:
offsets.append(offsets[-1] + state_size)
@@ -479,10 +886,12 @@ class GenericFluidSystem:
pneumatic_volume = self.pneumatic_volume_resolver.solve()
self.pneumatic_volume_propagation_count += pneumatic_volume.propagated
self._refresh_dynamic_components()
algebraic = self.pressure_flow_solver.solve(
initial_algebraic = self.pressure_flow_solver.solve(
effort_variables=("p",),
)
algebraic_diagnostics = [initial_algebraic]
pressure_flow_solve_count = 1
self.thermofluid_pressure_pass_count += 1
# Some constitutive flow laws recover their upstream temperature from
# connected stream enthalpy, while junction stream mixing itself depends
@@ -490,19 +899,25 @@ class GenericFluidSystem:
# leaves that two-way coupling to the next RHS call, making the ODE RHS
# depend on evaluation history and corrupting finite-difference
# Jacobians. Close both layers to one fixed point inside this call.
physical_ports = tuple(
port
for component in self.network.components.values()
for definition in component.active_port_definitions
if definition.kind == "physical"
for port in (component.get_port(definition.name),)
)
# The compiled closure plan keeps custom stream-aware components on the
# legacy global path. For catalog models, only stream-sensitive physical
# islands are revisited; independent islands keep the first solve.
closure_plan = self._thermofluid_closure_plan
self._last_algebraic_diagnostics = initial_algebraic
self._last_algebraic_scope = closure_plan.global_component_group
physical_ports = closure_plan.physical_ports
secondary_pressure_solvers = closure_plan.secondary_pressure_solvers
secondary_block_solvers = closure_plan.secondary_block_solvers
connected_h: dict[str, dict[str, float]] = {}
stream_diagnostics = []
max_coupling_iterations = 25
flow_relative_tolerance = 1.0e-12
for coupling_iteration in range(1, max_coupling_iterations + 1):
previous_flows = tuple(port.m_flow for port in physical_ports)
stream, connected_h = self.stream_resolver.solve()
stream, connected_h = self.stream_resolver.solve(
dynamic_ports_are_current=True,
)
stream_diagnostics.append(stream)
temperature_reference_h = (
self.stream_resolver.connected_temperature_reference_enthalpies()
)
@@ -511,12 +926,57 @@ class GenericFluidSystem:
component.update_flow_temperature_references(
temperature_reference_h[component.name]
)
algebraic = self.pressure_flow_solver.solve(
effort_variables=(
("p",) if pressure_flow_solve_count == 0 else ()
),
if secondary_pressure_solvers:
self.thermofluid_pressure_pass_count += 1
block_scale_context = (
self.pressure_flow_solver.scale_context()
if any(
solver is not self.pressure_flow_solver
for solver in secondary_pressure_solvers
)
else None
)
pressure_flow_solve_count += 1
if (
secondary_block_solvers
and not closure_plan.uses_conservative_global_solver
):
# Stream propagation only invalidates equations that explicitly
# consume the new enthalpy/temperature references. Re-solve the
# exact equation/unknown blocks containing those equations; the
# first global pass above remains the causalization boundary for
# mechanics, contact, and all stream-independent pneumatic blocks.
for block_solver in secondary_block_solvers:
block_result = block_solver.solve(
scale_context=block_scale_context,
)
# One public secondary closure is one logical solve. The
# block solver folds every local attempt and a possible
# accepted global fallback into this single diagnostic, so
# evaluations and blockFallbackUsed are counted exactly
# once here rather than once per internal equation block.
(algebraic,) = block_result.diagnostics
(scope,) = block_result.scopes
algebraic_diagnostics.append(algebraic)
self._last_algebraic_diagnostics = algebraic
self._last_algebraic_scope = scope
pressure_flow_solve_count += 1
else:
for pressure_solver, component_group in zip(
secondary_pressure_solvers,
closure_plan.secondary_component_groups,
):
algebraic = pressure_solver.solve(
effort_variables=(),
scale_context=(
block_scale_context
if pressure_solver is not self.pressure_flow_solver
else None
),
)
algebraic_diagnostics.append(algebraic)
self._last_algebraic_diagnostics = algebraic
self._last_algebraic_scope = component_group
pressure_flow_solve_count += 1
current_flows = tuple(port.m_flow for port in physical_ports)
flow_scale = max(
[abs(value) for value in (*previous_flows, *current_flows)] + [1.0]
@@ -528,7 +988,10 @@ class GenericFluidSystem:
),
default=0.0,
)
if max_flow_delta <= flow_relative_tolerance * flow_scale:
if (
not secondary_pressure_solvers
or max_flow_delta <= flow_relative_tolerance * flow_scale
):
break
else:
raise ThermofluidClosureError(
@@ -541,17 +1004,40 @@ class GenericFluidSystem:
)
self.mechanical_state_reducer.update_constraint_accelerations()
self.algebraic_solve_count += pressure_flow_solve_count
seeded_count = sum(
item.jacobian_mode == "seeded" for item in algebraic_diagnostics
)
self.algebraic_seeded_solve_count += seeded_count
self.algebraic_nonlinear_solve_count += (
len(algebraic_diagnostics) - seeded_count
)
self.algebraic_optimizer_evaluation_count += sum(
item.evaluations for item in algebraic_diagnostics
)
self.algebraic_residual_evaluation_count += sum(
item.residual_evaluations for item in algebraic_diagnostics
)
self.algebraic_block_fallback_count += sum(
item.block_fallback_used for item in algebraic_diagnostics
)
self.algebraic_dense_fallback_count += sum(
item.dense_fallback_used for item in algebraic_diagnostics
)
self.max_algebraic_residual = max(
self.max_algebraic_residual,
algebraic.max_scaled_residual,
*(item.max_scaled_residual for item in algebraic_diagnostics),
)
self.max_algebraic_evaluations = max(
self.max_algebraic_evaluations,
algebraic.evaluations,
*(item.evaluations for item in algebraic_diagnostics),
)
self.max_algebraic_residual_evaluations = max(
self.max_algebraic_residual_evaluations,
*(item.residual_evaluations for item in algebraic_diagnostics),
)
self.max_stream_iterations = max(
self.max_stream_iterations,
stream.iterations,
*(item.iterations for item in stream_diagnostics),
)
return connected_h
@@ -588,6 +1074,7 @@ class GenericFluidSystem:
f"{component.name}.{relative_key}", []
).append(value)
@with_property_cache
def simulate(
self,
config: SolveIVPConfig,
@@ -645,6 +1132,7 @@ class GenericFluidSystem:
report_progress(0.0, "integrating", force=True)
duration = config.t_stop - config.t_start
furthest_solver_time = config.t_start
next_signal_audit_index = 0
def report_solver_time(time: float) -> None:
nonlocal furthest_solver_time
@@ -657,10 +1145,23 @@ class GenericFluidSystem:
report_progress(time_fraction, "integrating")
def monitored_rhs(time: float, state_vector: list[float]) -> list[float]:
nonlocal next_signal_audit_index
while (
next_signal_audit_index < len(signal_event_times)
and float(time) >= signal_event_times[next_signal_audit_index]
):
self._request_causal_residual_audit()
next_signal_audit_index += 1
if cancel_check is None:
report_solver_time(time)
return self.rhs(time, state_vector)
def handle_state_transition(*args):
transition = self.mechanical_state_reducer.state_transition(*args)
if transition is not None:
self._request_causal_residual_audit()
return transition
solution = integrate_ode(
rhs=monitored_rhs,
initial_state=initial_state,
@@ -672,7 +1173,7 @@ class GenericFluidSystem:
),
breakpoints=signal_event_times,
state_transition_handler=(
self.mechanical_state_reducer.state_transition
handle_state_transition
if self.mechanical_state_reducer.has_state_events
else None
),
@@ -791,13 +1292,66 @@ class GenericFluidSystem:
},
"pressureFlow": {
"solveCount": self.algebraic_solve_count,
"seededSolveCount": self.algebraic_seeded_solve_count,
"nonlinearSolveCount": self.algebraic_nonlinear_solve_count,
"fastPathHitRate": (
self.algebraic_seeded_solve_count
/ self.algebraic_solve_count
if self.algebraic_solve_count
else 0.0
),
"optimizerEvaluationCount": (
self.algebraic_optimizer_evaluation_count
),
"residualEvaluationCount": (
self.algebraic_residual_evaluation_count
),
"blockFallbackCount": self.algebraic_block_fallback_count,
"denseFallbackCount": self.algebraic_dense_fallback_count,
"closurePassCount": self.thermofluid_pressure_pass_count,
"secondaryPhysicalIslandCount": len(
self._thermofluid_closure_plan.secondary_pressure_solvers
),
"secondaryBlockCount": sum(
len(solver.blocks)
for solver in self._thermofluid_closure_plan.secondary_block_solvers
if solver.available
),
"secondaryUnknownCount": sum(
len(block.unknowns)
for solver in self._thermofluid_closure_plan.secondary_block_solvers
if solver.available
for block in solver.blocks
),
"equationBlockFallbackReasons": [
solver.fallback_reason
for solver in self._thermofluid_closure_plan.secondary_block_solvers
if solver.fallback_reason is not None
],
"usesConservativeGlobalCoupling": (
self._thermofluid_closure_plan.uses_conservative_global_solver
),
"couplingPlanFallbackReason": (
self._thermofluid_closure_plan.conservative_fallback_reason
),
"maxScaledResidual": self.max_algebraic_residual,
"maxEvaluationsPerSolve": self.max_algebraic_evaluations,
"maxResidualEvaluationsPerSolve": (
self.max_algebraic_residual_evaluations
),
"lastScope": list(self._last_algebraic_scope),
"last": (
self.pressure_flow_solver.last_diagnostics.as_dict()
if self.pressure_flow_solver.last_diagnostics is not None
self._last_algebraic_diagnostics.as_dict()
if self._last_algebraic_diagnostics is not None
else None
),
"causalExecution": (
self.pressure_flow_solver.causal_execution_diagnostics()
),
"secondaryCausalExecution": [
solver.causal_execution_diagnostics()
for solver in self._thermofluid_closure_plan.secondary_block_solvers
],
},
"stream": {
"maxIterationsPerSolve": self.max_stream_iterations,
+161
View File
@@ -0,0 +1,161 @@
"""Process-local warm-up for the numerical simulation runtime."""
from __future__ import annotations
from dataclasses import asdict, dataclass
import logging
import os
from threading import Lock
from time import perf_counter
from typing import Literal
LOGGER = logging.getLogger(__name__)
WarmupStatus = Literal["completed", "failed", "disabled"]
@dataclass(frozen=True)
class SimulationWarmupReport:
status: WarmupStatus
duration_ms: float
error: str | None = None
def as_dict(self) -> dict[str, object]:
return asdict(self)
_WARMUP_LOCK = Lock()
_WARMUP_REPORT: SimulationWarmupReport | None = None
def simulation_warmup_enabled() -> bool:
raw_value = os.getenv("SIMULATIONAPP_WARMUP", "on").strip().lower()
if raw_value in {"", "1", "true", "yes", "on"}:
return True
if raw_value in {"0", "false", "no", "off"}:
return False
raise ValueError(
"SIMULATIONAPP_WARMUP must be one of: on, off, true, false, 1, 0."
)
def _run_numerical_warmup() -> None:
"""Exercise only in-memory SciPy paths used by real simulations."""
import numpy as np
from scipy.integrate import BDF, DOP853, LSODA, RK23, RK45, Radau, solve_ivp
from scipy.optimize import brentq, least_squares
from scipy.optimize._numdiff import group_columns
from scipy.sparse import csc_matrix, csr_matrix
# Importing these classes is intentional even though the micro solve below
# uses BDF: the stepwise solver selects them dynamically at runtime.
solver_types = (BDF, DOP853, LSODA, RK23, RK45, Radau)
if len(solver_types) != 6:
raise RuntimeError("SciPy solver warm-up did not load every supported method.")
sparsity = csc_matrix(np.array([[1.0]], dtype=float))
groups = group_columns(sparsity)
if groups.shape != (1,):
raise RuntimeError("SciPy Jacobian grouping warm-up returned an invalid shape.")
integration = solve_ivp(
lambda _time, state: -state,
(0.0, 1.0e-4),
np.array([1.0], dtype=float),
method="BDF",
t_eval=np.array([0.0, 1.0e-4], dtype=float),
jac_sparsity=sparsity,
rtol=1.0e-6,
atol=1.0e-9,
)
if not integration.success or not np.isfinite(integration.y).all():
raise RuntimeError("SciPy integration warm-up did not complete successfully.")
algebraic_sparsity = csr_matrix(np.eye(2, dtype=bool))
algebraic = least_squares(
lambda state: np.array(
[state[0] - 1.0, state[1] - 2.0],
dtype=float,
),
np.array([0.5, 0.5], dtype=float),
bounds=(
np.array([0.0, 0.0], dtype=float),
np.array([3.0, 3.0], dtype=float),
),
jac_sparsity=algebraic_sparsity,
tr_solver="lsmr",
)
if (
not algebraic.success
or not np.isfinite(algebraic.x).all()
or not np.allclose(algebraic.x, np.array([1.0, 2.0]), atol=1.0e-8)
):
raise RuntimeError("SciPy algebraic warm-up did not complete successfully.")
root = brentq(lambda value: value - 0.5, 0.0, 1.0)
if abs(root - 0.5) > 1.0e-12:
raise RuntimeError("SciPy scalar root warm-up returned an invalid result.")
# Compile the cached v3 XSD through the same public validation path. The
# intentionally incomplete document is never accepted or persisted.
from app.system_xml import validate_system_xml_document
validate_system_xml_document(b"<System/>")
def warm_up_simulation_runtime() -> SimulationWarmupReport:
"""Warm one worker exactly once, returning a startup diagnostic report.
Ordinary warm-up failures are reported but do not prevent the editor and
non-simulation APIs from starting. ``MemoryError`` remains fatal because
continuing a worker under memory exhaustion is unsafe.
"""
global _WARMUP_REPORT
with _WARMUP_LOCK:
if _WARMUP_REPORT is not None:
return _WARMUP_REPORT
if not simulation_warmup_enabled():
_WARMUP_REPORT = SimulationWarmupReport(
status="disabled",
duration_ms=0.0,
)
return _WARMUP_REPORT
started = perf_counter()
try:
_run_numerical_warmup()
except MemoryError:
raise
except Exception as exc:
_WARMUP_REPORT = SimulationWarmupReport(
status="failed",
duration_ms=(perf_counter() - started) * 1000.0,
error=f"{type(exc).__name__}: {exc}",
)
LOGGER.exception("Simulation runtime warm-up failed; startup will continue.")
else:
_WARMUP_REPORT = SimulationWarmupReport(
status="completed",
duration_ms=(perf_counter() - started) * 1000.0,
)
LOGGER.info(
"Simulation runtime warm-up completed in %.1f ms.",
_WARMUP_REPORT.duration_ms,
)
return _WARMUP_REPORT
def _reset_simulation_warmup_for_tests() -> None:
global _WARMUP_REPORT
with _WARMUP_LOCK:
_WARMUP_REPORT = None
__all__ = [
"SimulationWarmupReport",
"simulation_warmup_enabled",
"warm_up_simulation_runtime",
]