完善通用求解器回归与前端交互

- 引入因果坐标内核、热流体恢复和递进长时回归\n- 完善正交连线、线桥、视图保持与结果曲线缩放\n- 补充依赖约束、CI、测试基线和北京时间更新日志
This commit is contained in:
lujingze committed 2026-08-18 06:42:07 +00:00
1 parent 143e8dd309
commit b435daecf2
65 files changed
+172271 -701

No files matched your search

+492 -26
View File
@@ -31,6 +31,10 @@ PRESSURE_LOWER_BOUND_PA = 0.0
CAUSAL_FAST_PATH_ENVIRONMENT_VARIABLE = "SIMULATION_CAUSAL_FAST_PATH"
CAUSAL_EXECUTOR_V2_ENVIRONMENT_VARIABLE = "SIMULATION_CAUSAL_EXECUTOR_V2"
CAUSAL_COORDINATE_KERNEL_ENVIRONMENT_VARIABLE = (
"SIMULATION_CAUSAL_COORDINATE_KERNEL"
)
CAUSAL_FAST_PATH_AUDIT_INTERVAL = 64
@@ -39,6 +43,20 @@ def _causal_fast_path_environment_enabled() -> bool:
return value.strip().lower() not in {"0", "false", "no", "off"}
def _causal_executor_v2_environment_enabled() -> bool:
"""Return whether the allocation-light causal executor is enabled."""
value = os.getenv(CAUSAL_EXECUTOR_V2_ENVIRONMENT_VARIABLE, "1")
return value.strip().lower() not in {"0", "false", "no", "off"}
def _causal_coordinate_kernel_environment_enabled() -> bool:
"""Return whether the canonical-coordinate causal kernel is enabled."""
value = os.getenv(CAUSAL_COORDINATE_KERNEL_ENVIRONMENT_VARIABLE, "1")
return value.strip().lower() not in {"0", "false", "no", "off"}
class AlgebraicSolveError(RuntimeError):
def __init__(
self,
@@ -120,6 +138,42 @@ class CausalEffortAssignment:
anchor: EffortAnchor
@dataclass(frozen=True)
class CausalEffortKernelTarget:
"""One canonical effort coordinate extracted from a residual evaluator."""
coordinate_index: int
assignment: CausalEffortAssignment
equation_index: int
equation_id: str
@dataclass(frozen=True)
class CausalEffortKernelEvaluation:
"""One component call shared by every state anchor that it owns."""
evaluate: Callable[[], tuple[float, ...]]
targets: tuple[CausalEffortKernelTarget, ...]
@dataclass(frozen=True)
class CausalEffortKernelStage:
"""Precompiled canonical coordinates and compatibility broadcasts."""
variable: str
assignments: tuple[tuple[int, CausalEffortAssignment], ...]
direct_targets: tuple[tuple[int, CausalEffortAssignment], ...]
component_evaluations: tuple[CausalEffortKernelEvaluation, ...]
@dataclass(frozen=True)
class CausalFlowKernelStage:
"""Map one existing dependency stage into the canonical workspace."""
stage: ExplicitFlowStage
coordinate_indices: tuple[int, ...]
@dataclass(frozen=True)
class ConnectionEquationEvaluation:
template: EquationResidual
@@ -366,6 +420,56 @@ class PressureFlowSolver:
self._causal_fast_path_environment_enabled = (
_causal_fast_path_environment_enabled()
)
self._causal_executor_v2_environment_enabled = (
_causal_executor_v2_environment_enabled()
)
self._causal_coordinate_kernel_environment_enabled = (
_causal_coordinate_kernel_environment_enabled()
)
self._causal_compiled_effort_unknown_count = sum(
len(assignment.members)
for assignments in self._causal_effort_plan_by_variable.values()
for assignment in assignments
)
self._causal_compiled_flow_assignment_count = sum(
len(stage.assignments) for stage in self._explicit_flow_plan
)
self._causal_external_effort_unknowns = tuple(
member
for variable in ("x", "v")
for assignment in self._causal_effort_plan_by_variable.get(
variable, ()
)
for member in assignment.members
)
self._causal_external_x_states = tuple(
unknown.state
for unknown in self._causal_external_effort_unknowns
if unknown.variable == "x"
)
self._causal_external_v_states = tuple(
unknown.state
for unknown in self._causal_external_effort_unknowns
if unknown.variable == "v"
)
(
self._causal_effort_kernel_by_variable,
self._causal_flow_kernel_plan,
self._causal_coordinate_values,
) = self._compile_causal_coordinate_kernel()
self._causal_logical_effort_coordinate_count = sum(
len(stage.assignments)
for stage in self._causal_effort_kernel_by_variable.values()
)
self._causal_eliminated_effort_alias_count = max(
self._causal_compiled_effort_unknown_count
- self._causal_logical_effort_coordinate_count,
0,
)
self._causal_compatibility_scatter_count = (
self._causal_compiled_effort_unknown_count
+ self._causal_compiled_flow_assignment_count
)
self._causal_runtime_disabled_reason: str | None = None
self._causal_audit_interval = CAUSAL_FAST_PATH_AUDIT_INTERVAL
self._causal_audit_required = True
@@ -374,7 +478,11 @@ class PressureFlowSolver:
self._causal_full_residual_audit_count = 0
self._causal_audit_failure_count = 0
self._causal_legacy_fallback_count = 0
self._causal_v2_fast_solve_count = 0
self._causal_v2_runtime_validation_failure_count = 0
self._causal_coordinate_fast_solve_count = 0
self._causal_last_verified_diagnostics: AlgebraicSolveDiagnostics | None = None
self._causal_cached_fast_diagnostics: AlgebraicSolveDiagnostics | None = None
self.last_diagnostics: AlgebraicSolveDiagnostics | None = None
@property
@@ -435,6 +543,25 @@ class PressureFlowSolver:
and self._causal_runtime_disabled_reason is None
)
@property
def causal_executor_v2_enabled(self) -> bool:
"""Whether this run may use the v2 executor (environment opt-out)."""
return (
self._causal_executor_v2_environment_enabled
and self.causal_fast_path_enabled
)
@property
def causal_coordinate_kernel_enabled(self) -> bool:
"""Whether the canonical-coordinate executor may run now."""
return (
self._causal_coordinate_kernel_environment_enabled
and self.causal_executor_v2_enabled
and bool(self._causal_coordinate_values)
)
def causal_execution_diagnostics(self) -> dict[str, object]:
disabled_reason = self._causal_runtime_disabled_reason
if not self._causal_fast_path_environment_enabled:
@@ -458,6 +585,39 @@ class PressureFlowSolver:
if last_verified is not None
else None
),
"executorV2Configured": self._causal_executor_v2_environment_enabled,
"executorV2Enabled": self.causal_executor_v2_enabled,
"executorV2FastSolveCount": self._causal_v2_fast_solve_count,
"executorV2RuntimeValidationFailureCount": (
self._causal_v2_runtime_validation_failure_count
),
"coordinateKernelConfigured": (
self._causal_coordinate_kernel_environment_enabled
),
"coordinateKernelEnabled": self.causal_coordinate_kernel_enabled,
"coordinateKernelFastSolveCount": (
self._causal_coordinate_fast_solve_count
),
"compiledEffortUnknownCount": (
self._causal_compiled_effort_unknown_count
),
"compiledFlowAssignmentCount": (
self._causal_compiled_flow_assignment_count
),
"compiledAssignmentCount": (
self._causal_compiled_effort_unknown_count
+ self._causal_compiled_flow_assignment_count
),
"logicalEffortCoordinateCount": (
self._causal_logical_effort_coordinate_count
),
"eliminatedEffortAliasCount": (
self._causal_eliminated_effort_alias_count
),
"canonicalCoordinateCount": len(self._causal_coordinate_values),
"compatibilityScatterCount": (
self._causal_compatibility_scatter_count
),
}
def request_causal_audit(self) -> None:
@@ -651,10 +811,163 @@ class PressureFlowSolver:
None,
)
def _compile_causal_coordinate_kernel(
self,
) -> tuple[
dict[str, CausalEffortKernelStage],
tuple[CausalFlowKernelStage, ...],
list[float],
]:
"""Compile independent coordinates without changing public port state.
``PortState`` remains the compatibility surface consumed by component
methods. The workspace stores one value per proven effort equality
group and one per explicit flow assignment; compatibility aliases are
populated only after every target in an effort stage has been checked.
"""
if not self._causal_fast_path_eligible:
return {}, (), []
equation_index_by_id = {
equation.id: index
for index, equation in enumerate(self._equation_templates)
}
effort_stages: dict[str, CausalEffortKernelStage] = {}
next_coordinate = 0
for variable in ("p", "x", "v"):
assignments = self._causal_effort_plan_by_variable.get(variable, ())
indexed_assignments = tuple(
(next_coordinate + offset, assignment)
for offset, assignment in enumerate(assignments)
)
next_coordinate += len(indexed_assignments)
direct_targets: list[tuple[int, CausalEffortAssignment]] = []
targets_by_component: dict[
int,
list[CausalEffortKernelTarget],
] = {}
component_evaluators: dict[int, Callable[[], tuple[float, ...]]] = {}
for coordinate_index, assignment in indexed_assignments:
equation_index = equation_index_by_id[
assignment.anchor.equation_id
]
kind, evaluation_plan, source = (
self._equation_evaluation_locations[equation_index]
)
if kind != "component":
direct_targets.append((coordinate_index, assignment))
continue
component_plan = evaluation_plan
key = id(component_plan)
component_evaluators[key] = component_plan.evaluate
targets_by_component.setdefault(key, []).append(
CausalEffortKernelTarget(
coordinate_index=coordinate_index,
assignment=assignment,
equation_index=source,
equation_id=assignment.anchor.equation_id,
)
)
effort_stages[variable] = CausalEffortKernelStage(
variable=variable,
assignments=indexed_assignments,
direct_targets=tuple(direct_targets),
component_evaluations=tuple(
CausalEffortKernelEvaluation(
evaluate=component_evaluators[key],
targets=tuple(targets),
)
for key, targets in targets_by_component.items()
),
)
flow_stages: list[CausalFlowKernelStage] = []
for stage in self._explicit_flow_plan:
coordinate_indices = tuple(
range(next_coordinate, next_coordinate + len(stage.assignments))
)
next_coordinate += len(stage.assignments)
flow_stages.append(
CausalFlowKernelStage(
stage=stage,
coordinate_indices=coordinate_indices,
)
)
return effort_stages, tuple(flow_stages), [0.0] * next_coordinate
@staticmethod
def _read_effort_anchor(assignment: CausalEffortAssignment) -> float:
state = assignment.anchor.unknown.state
if assignment.variable == "p":
return state.p
if assignment.variable == "x":
return state.x
return state.v
@staticmethod
def _scatter_effort_assignment(
assignment: CausalEffortAssignment,
value: float,
) -> None:
if assignment.variable == "p":
for unknown in assignment.members:
unknown.state.p = value
return
if assignment.variable == "x":
for unknown in assignment.members:
unknown.state.x = value
return
for unknown in assignment.members:
unknown.state.v = value
def _execute_causal_coordinate_effort_plan(
self,
variables: tuple[str, ...],
) -> bool:
"""Evaluate canonical effort coordinates in component-sized batches."""
workspace = self._causal_coordinate_values
for variable in variables:
stage = self._causal_effort_kernel_by_variable.get(variable)
if stage is None:
return False
for coordinate_index, assignment in stage.direct_targets:
workspace[coordinate_index] = (
self._read_effort_anchor(assignment)
- assignment.anchor.evaluate()
)
for evaluation in stage.component_evaluations:
equation_values = evaluation.evaluate()
for target in evaluation.targets:
if target.equation_index >= len(equation_values):
raise RuntimeError(
"Compiled algebraic equation disappeared at runtime: "
f"{target.equation_id}."
)
workspace[target.coordinate_index] = (
self._read_effort_anchor(target.assignment)
- float(equation_values[target.equation_index])
)
for coordinate_index, assignment in stage.assignments:
target = workspace[coordinate_index]
if not isfinite(target) or (
variable == "p" and target <= PRESSURE_LOWER_BOUND_PA
):
return False
for coordinate_index, assignment in stage.assignments:
self._scatter_effort_assignment(
assignment,
workspace[coordinate_index],
)
return True
def _execute_causal_effort_plan(
self,
variables: tuple[str, ...],
) -> bool:
if self.causal_coordinate_kernel_enabled:
return self._execute_causal_coordinate_effort_plan(variables)
for variable in variables:
assignments = self._causal_effort_plan_by_variable.get(variable)
if assignments is None:
@@ -685,15 +998,8 @@ class PressureFlowSolver:
self._causal_solves_since_audit = 0
self._causal_audit_required = False
self._causal_last_verified_diagnostics = diagnostics
def _causal_fast_diagnostics(
self,
) -> AlgebraicSolveDiagnostics:
verified = self._causal_last_verified_diagnostics
if verified is None:
raise RuntimeError("Causal execution has no verified residual baseline.")
return replace(
verified,
self._causal_cached_fast_diagnostics = replace(
diagnostics,
message=(
"Compiled causal pressure-flow program completed; residuals "
"reuse the latest full audit."
@@ -709,6 +1015,14 @@ class PressureFlowSolver:
causal_fast_path_used=True,
)
def _causal_fast_diagnostics(
self,
) -> AlgebraicSolveDiagnostics:
cached = self._causal_cached_fast_diagnostics
if cached is None:
raise RuntimeError("Causal execution has no verified residual baseline.")
return cached
def _build_jacobian_sparsity(self):
"""Compile the residual dependency contract into one CSR pattern.
@@ -1024,6 +1338,10 @@ class PressureFlowSolver:
unknown = sorted(set(variables) - set(self._effort_groups))
if unknown:
raise ValueError("Unsupported effort variables: " + ", ".join(unknown))
if self.causal_coordinate_kernel_enabled:
if self._execute_causal_coordinate_effort_plan(variables):
return
self._disable_causal_fast_path("nonFiniteCausalEffortAnchor")
for variable in variables:
self._seed_equal_effort(variable)
@@ -1835,6 +2153,121 @@ class PressureFlowSolver:
seeded_ids.add(assignment.unknown.id)
return seeded_ids
def _execute_compiled_causal_flow_plan(self) -> str | None:
"""Execute the compile-proven full flow plan without coverage sets."""
if self.causal_coordinate_kernel_enabled:
return self._execute_causal_coordinate_flow_plan()
# Position and velocity are propagated by the mechanical reducer
# before the pressure-only causal solve. They are therefore not
# rewritten below, but remain part of the compiled algebraic contract.
# Validate that small external boundary explicitly instead of restoring
# the legacy scan over every pressure/flow/force unknown.
if any(
not isfinite(unknown.read())
for unknown in self._causal_external_effort_unknowns
):
return "nonFiniteCausalExternalEffort"
reset_unknowns = self._explicit_flow_unknowns_by_variables[
frozenset(("f", "m_flow"))
]
for unknown in reset_unknowns:
unknown.write(0.0)
for stage in self._explicit_flow_plan:
try:
values = self._evaluate_explicit_flow_stage(stage)
except MemoryError:
raise
except (ArithmeticError, RuntimeError, ValueError) as exc:
return f"causalFlowEvaluationFailed:{type(exc).__name__}"
if len(values) != len(stage.assignments):
return "causalFlowAssignmentCountMismatch"
for assignment, target_value in zip(stage.assignments, values):
if not isfinite(target_value):
return "nonFiniteCausalFlowAssignment"
assignment.unknown.write(target_value)
return None
def _execute_causal_coordinate_flow_plan(self) -> str | None:
"""Run flow stages through reusable canonical coordinates."""
if any(not isfinite(state.x) for state in self._causal_external_x_states):
return "nonFiniteCausalExternalEffort"
if any(not isfinite(state.v) for state in self._causal_external_v_states):
return "nonFiniteCausalExternalEffort"
return self._execute_causal_coordinate_flow_stages(
self._causal_flow_kernel_plan,
self._causal_coordinate_values,
)
@staticmethod
def _execute_causal_coordinate_flow_stages(
kernel_plan: tuple[CausalFlowKernelStage, ...],
workspace: list[float],
) -> str | None:
"""Execute proven flow stages without per-call result containers."""
# Residual-based explicit assignments use ``-residual`` and therefore
# require their target coordinate to be zero. Keep this compatibility
# initialization until a component exposes a proven direct target op.
for kernel_stage in kernel_plan:
for assignment in kernel_stage.stage.assignments:
if assignment.unknown.variable == "m_flow":
assignment.unknown.state.m_flow = 0.0
else:
assignment.unknown.state.f = 0.0
for kernel_stage in kernel_plan:
stage = kernel_stage.stage
coordinate_indices = kernel_stage.coordinate_indices
if len(coordinate_indices) != len(stage.assignments):
return "causalFlowAssignmentCountMismatch"
try:
for assignment_index, evaluate in stage.direct_evaluations:
workspace[coordinate_indices[assignment_index]] = float(
evaluate()
)
for evaluation in stage.component_evaluations:
equation_values = evaluation.evaluate()
for assignment_index, equation_index, equation_id in zip(
evaluation.assignment_indices,
evaluation.equation_indices,
evaluation.equation_ids,
):
if equation_index >= len(equation_values):
raise RuntimeError(
"Compiled algebraic equation disappeared at "
f"runtime: {equation_id}."
)
workspace[coordinate_indices[assignment_index]] = (
0.0 - float(equation_values[equation_index])
)
except MemoryError:
raise
except (ArithmeticError, RuntimeError, ValueError) as exc:
return f"causalFlowEvaluationFailed:{type(exc).__name__}"
for assignment, coordinate_index in zip(
stage.assignments,
coordinate_indices,
):
target_value = workspace[coordinate_index]
if not isfinite(target_value):
return "nonFiniteCausalFlowAssignment"
for assignment, coordinate_index in zip(
stage.assignments,
coordinate_indices,
):
target_value = workspace[coordinate_index]
if assignment.unknown.variable == "m_flow":
assignment.unknown.state.m_flow = target_value
else:
assignment.unknown.state.f = target_value
return None
def _build_closed_resistance_pressure_plan(
self,
) -> tuple[ClosedResistancePressureBinding, ...]:
@@ -2088,6 +2521,9 @@ class PressureFlowSolver:
causal_audit_due = (
self._causal_audit_is_due() if causal_candidate else False
)
causal_v2_candidate = (
causal_candidate and self._causal_executor_v2_environment_enabled
)
for component in self._causal_contact_components:
component.clear_causal_contact()
@@ -2100,29 +2536,55 @@ class PressureFlowSolver:
self._disable_causal_fast_path("nonFiniteCausalEffortAnchor")
causal_candidate = False
causal_audit_due = False
causal_v2_candidate = False
self._seed_equal_efforts(effort_variables)
else:
self._seed_equal_efforts(effort_variables)
self._seed_closed_resistance_pressures()
self._seed_resistance_pnl0001_series_pressures()
seeded_flow_ids = self._solve_explicit_flow_unknowns()
contact_bindings = self._seed_unilateral_contacts()
if contact_bindings:
seeded_flow_ids.update(self._solve_explicit_flow_unknowns(("f",)))
self._refresh_unilateral_contacts(contact_bindings)
seeded_flow_ids: set[str] | None = None
contact_bindings: tuple[UnilateralContactBinding, ...] = ()
if causal_v2_candidate:
v2_failure_reason = self._execute_compiled_causal_flow_plan()
if v2_failure_reason is not None:
self._causal_v2_runtime_validation_failure_count += 1
self._causal_legacy_fallback_count += 1
self._disable_causal_fast_path(v2_failure_reason)
causal_candidate = False
causal_audit_due = False
causal_v2_candidate = False
# Rebuild the ordinary seed from scratch in the same solve.
# A partial compiled stage must never influence fallback.
self._seed_equal_efforts(effort_variables)
if not causal_v2_candidate:
self._seed_closed_resistance_pressures()
self._seed_resistance_pnl0001_series_pressures()
seeded_flow_ids = self._solve_explicit_flow_unknowns()
contact_bindings = self._seed_unilateral_contacts()
if contact_bindings:
seeded_flow_ids.update(
self._solve_explicit_flow_unknowns(("f",))
)
self._refresh_unilateral_contacts(contact_bindings)
if causal_candidate:
causal_unknowns_are_feasible = (
not contact_bindings
and seeded_flow_ids == self._causal_flow_unknown_ids
and all(
isfinite(unknown.read())
and (
unknown.variable != "p"
or unknown.read() > PRESSURE_LOWER_BOUND_PA
if causal_v2_candidate:
# Compilation proves a disjoint, complete effort/flow
# partition. The v2 executors validate each produced value,
# so no coverage set or full unknown scan is needed here.
causal_unknowns_are_feasible = True
else:
assert seeded_flow_ids is not None
causal_unknowns_are_feasible = (
not contact_bindings
and seeded_flow_ids == self._causal_flow_unknown_ids
and all(
isfinite(unknown.read())
and (
unknown.variable != "p"
or unknown.read() > PRESSURE_LOWER_BOUND_PA
)
for unknown in self.unknowns
)
for unknown in self.unknowns
)
)
if not causal_unknowns_are_feasible:
self._causal_legacy_fallback_count += 1
self._disable_causal_fast_path("causalRuntimeGateFailed")
@@ -2131,6 +2593,10 @@ class PressureFlowSolver:
elif not causal_audit_due:
diagnostics = self._causal_fast_diagnostics()
self._causal_fast_solve_count += 1
if causal_v2_candidate:
self._causal_v2_fast_solve_count += 1
if self.causal_coordinate_kernel_enabled:
self._causal_coordinate_fast_solve_count += 1
self._causal_solves_since_audit += 1
self.last_diagnostics = diagnostics
return diagnostics
+177 -2
View File
@@ -9,6 +9,7 @@ from app.simulation.solvers.algebraic import (
PRESSURE_LOWER_BOUND_PA,
AlgebraicSolveDiagnostics,
AlgebraicUnknown,
CausalFlowKernelStage,
ExplicitFlowStage,
PressureFlowSolver,
)
@@ -237,6 +238,27 @@ class StreamPressureBlockSolver:
)
for stage in pressure_flow_solver._explicit_flow_plan
)
secondary_coordinate = 0
secondary_kernel_plan: list[CausalFlowKernelStage] = []
for stage in self._selected_explicit_flow_plan:
coordinate_indices = tuple(
range(
secondary_coordinate,
secondary_coordinate + len(stage.assignments),
)
)
secondary_coordinate += len(stage.assignments)
secondary_kernel_plan.append(
CausalFlowKernelStage(
stage=stage,
coordinate_indices=coordinate_indices,
)
)
self._causal_secondary_flow_kernel_plan = tuple(secondary_kernel_plan)
self._causal_secondary_coordinate_values = [0.0] * secondary_coordinate
self._causal_v2_entry_values = [0.0] * len(
self._selected_flow_unknowns
)
self._selected_equation_evaluation = (
self._compile_selected_equation_evaluation()
if self.blocks
@@ -259,7 +281,11 @@ class StreamPressureBlockSolver:
self._causal_full_residual_audit_count = 0
self._causal_audit_failure_count = 0
self._causal_legacy_fallback_count = 0
self._causal_v2_fast_solve_count = 0
self._causal_v2_runtime_validation_failure_count = 0
self._causal_coordinate_fast_solve_count = 0
self._causal_last_verified_diagnostics: AlgebraicSolveDiagnostics | None = None
self._causal_cached_fast_diagnostics: AlgebraicSolveDiagnostics | None = None
@property
def available(self) -> bool:
@@ -273,6 +299,13 @@ class StreamPressureBlockSolver:
and self.pressure_flow_solver.causal_fast_path_enabled
)
@property
def causal_executor_v2_enabled(self) -> bool:
return (
self.causal_fast_path_enabled
and self.pressure_flow_solver._causal_executor_v2_environment_enabled
)
def request_causal_audit(self) -> None:
self._causal_audit_required = True
@@ -280,7 +313,9 @@ class StreamPressureBlockSolver:
parent = self.pressure_flow_solver.causal_execution_diagnostics()
disabled_reason = self._causal_runtime_disabled_reason
if not bool(parent["enabled"]):
disabled_reason = str(parent["disabledReason"] or "parentCausalPathDisabled")
disabled_reason = str(
parent["disabledReason"] or "parentCausalPathDisabled"
)
elif not self._causal_fast_path_eligible:
disabled_reason = self._causal_fast_path_fallback_reason
verified = self._causal_last_verified_diagnostics
@@ -298,6 +333,34 @@ class StreamPressureBlockSolver:
"lastVerifiedMaxScaledResidual": (
verified.max_scaled_residual if verified is not None else None
),
"executorV2Configured": (
self.pressure_flow_solver._causal_executor_v2_environment_enabled
),
"executorV2Enabled": self.causal_executor_v2_enabled,
"executorV2FastSolveCount": self._causal_v2_fast_solve_count,
"executorV2RuntimeValidationFailureCount": (
self._causal_v2_runtime_validation_failure_count
),
"coordinateKernelConfigured": (
self.pressure_flow_solver._causal_coordinate_kernel_environment_enabled
),
"coordinateKernelEnabled": (
self.causal_executor_v2_enabled
and self.pressure_flow_solver.causal_coordinate_kernel_enabled
),
"coordinateKernelFastSolveCount": (
self._causal_coordinate_fast_solve_count
),
"compiledEffortUnknownCount": len(
self._causal_effort_entry_positions
),
"compiledFlowAssignmentCount": len(
self._causal_expected_flow_equation_ids
),
"compiledAssignmentCount": (
len(self._causal_effort_entry_positions)
+ len(self._causal_expected_flow_equation_ids)
),
}
def _disable_causal_fast_path(self, reason: str) -> None:
@@ -405,6 +468,28 @@ class StreamPressureBlockSolver:
self._causal_solves_since_audit = 0
self._causal_audit_required = False
self._causal_last_verified_diagnostics = diagnostics
self._causal_cached_fast_diagnostics = replace(
diagnostics,
message=(
"Compiled causal stream-pressure block completed; residuals "
"reuse the latest full audit."
),
evaluations=0,
residual_evaluations=0,
dense_fallback_used=False,
nonlinear_block_count=0,
nonlinear_block_unknown_count=0,
block_fallback_used=False,
block_fallback_reason=None,
residual_verified_this_solve=False,
causal_fast_path_used=True,
)
def _causal_v2_fast_diagnostics(self) -> AlgebraicSolveDiagnostics:
cached = self._causal_cached_fast_diagnostics
if cached is None:
raise RuntimeError("Stream causal execution has no residual audit.")
return cached
def _causal_fast_diagnostics(
self,
@@ -788,6 +873,45 @@ class StreamPressureBlockSolver:
unknown.write(entry_values[position])
return frozenset(seeded_equation_ids)
def _execute_compiled_secondary_flow_plan(self) -> str | None:
"""Execute selected flow assignments without equation-id sets."""
if self.pressure_flow_solver.causal_coordinate_kernel_enabled:
failure = (
self.pressure_flow_solver._execute_causal_coordinate_flow_stages(
self._causal_secondary_flow_kernel_plan,
self._causal_secondary_coordinate_values,
)
)
if failure == "causalFlowAssignmentCountMismatch":
return "causalSecondaryFlowAssignmentCountMismatch"
if failure == "nonFiniteCausalFlowAssignment":
return "nonFiniteCausalSecondaryFlowAssignment"
if failure and failure.startswith("causalFlowEvaluationFailed:"):
return "causalSecondaryFlowEvaluationFailed:" + failure.rsplit(
":", 1
)[-1]
return failure
for unknown in self._selected_flow_unknowns:
unknown.write(0.0)
for stage in self._selected_explicit_flow_plan:
try:
values = self.pressure_flow_solver._evaluate_explicit_flow_stage(
stage
)
except MemoryError:
raise
except (ArithmeticError, RuntimeError, ValueError) as exc:
return f"causalSecondaryFlowEvaluationFailed:{type(exc).__name__}"
if len(values) != len(stage.assignments):
return "causalSecondaryFlowAssignmentCountMismatch"
for assignment, target_value in zip(stage.assignments, values):
if not isfinite(target_value):
return "nonFiniteCausalSecondaryFlowAssignment"
assignment.unknown.write(target_value)
return None
@staticmethod
def _equation_scales_from_specs(
specs: tuple[_EquationScaleSpec, ...],
@@ -1083,11 +1207,62 @@ class StreamPressureBlockSolver:
scale_context: Mapping[str, float] | None = None,
) -> StreamBlockSolveResult:
solver = self.pressure_flow_solver
context = dict(scale_context or solver.scale_context())
causal_candidate = self.causal_fast_path_enabled
causal_audit_due = (
self._causal_audit_is_due() if causal_candidate else False
)
causal_v2_candidate = (
causal_candidate
and solver._causal_executor_v2_environment_enabled
and not causal_audit_due
)
if causal_v2_candidate:
# The secondary causal proof rejects every special pressure seed,
# so this executor mutates selected flow coordinates only. Keep
# the minimal transactional snapshot for the rare fallback path.
v2_entry_values = self._causal_v2_entry_values
for position, unknown in enumerate(self._selected_flow_unknowns):
v2_entry_values[position] = unknown.state.m_flow
def restore_v2_entry_mutations() -> None:
for unknown, value in zip(
self._selected_flow_unknowns,
v2_entry_values,
):
unknown.state.m_flow = value
try:
v2_failure_reason = (
self._execute_compiled_secondary_flow_plan()
)
except BaseException:
restore_v2_entry_mutations()
raise
if v2_failure_reason is None:
diagnostics = self._causal_v2_fast_diagnostics()
self._causal_fast_solve_count += 1
self._causal_v2_fast_solve_count += 1
if solver.causal_coordinate_kernel_enabled:
self._causal_coordinate_fast_solve_count += 1
self._causal_solves_since_audit += 1
selected = self._selected_equation_evaluation
assert selected is not None
return StreamBlockSolveResult(
diagnostics=(diagnostics,),
scopes=(selected.scope_components,),
used_global_fallback=False,
)
restore_v2_entry_mutations()
self._causal_v2_runtime_validation_failure_count += 1
self._causal_legacy_fallback_count += 1
self._disable_causal_fast_path(v2_failure_reason)
causal_candidate = False
causal_audit_due = False
# Keep scale construction and the full mutation snapshot off the v2
# success path. Callers may still precompute a shared scale mapping;
# avoiding that producer requires a later Generic-system API change.
context = dict(scale_context or solver.scale_context())
entry_values = tuple(
unknown.read() for unknown in self._entry_mutated_unknowns
)
+783
View File
@@ -0,0 +1,783 @@
"""Executable reference IR for compile-proven causal algebraic programs.
The IR eliminates duplicate *logical* effort coordinates, but intentionally
keeps a compatibility scatter map to existing ``PortState`` objects. Stream
propagation, derivatives, and result collection still consume those objects;
this is a reference for a future flat backend, not physical slot deletion.
"""
from __future__ import annotations
from collections.abc import Callable, Iterable
from dataclasses import dataclass, replace
from enum import StrEnum
from hashlib import sha256
import json
from math import isfinite
from typing import TYPE_CHECKING, Any
if TYPE_CHECKING:
import numpy as np
CAUSAL_NUMERIC_IR_SCHEMA_VERSION = 1
PRESSURE_LOWER_BOUND_PA = 0.0
class CausalIROpcode(StrEnum):
EFFORT_BROADCAST = "effort_broadcast"
EFFORT_DIRECT_RESIDUAL = "effort_direct_residual"
EFFORT_COMPONENT_RESIDUAL = "effort_component_residual"
FLOW_DIRECT = "flow_direct"
FLOW_COMPONENT_RESIDUAL = "flow_component_residual"
@dataclass(frozen=True, slots=True)
class CausalIRCompatibilitySlot:
slot: int
id: str
variable: str
@dataclass(frozen=True, slots=True)
class CausalIRCanonicalSlot:
slot: int
id: str
variable: str
kind: str
@dataclass(frozen=True, slots=True)
class CausalIREffortOperation:
opcode: CausalIROpcode
variable: str
result_slot: int
anchor_compatibility_slot: int
scatter_compatibility_slots: tuple[int, ...]
equation_id: str
@dataclass(frozen=True, slots=True)
class CausalIREffortEvaluation:
opcode: CausalIROpcode
output_indices: tuple[int, ...]
equation_indices: tuple[int, ...]
equation_ids: tuple[str, ...]
evaluator_slot: int
@dataclass(frozen=True, slots=True)
class CausalIREffortStage:
variable: str
operations: tuple[CausalIREffortOperation, ...]
evaluations: tuple[CausalIREffortEvaluation, ...]
@dataclass(frozen=True, slots=True)
class CausalIRFlowOperation:
opcode: CausalIROpcode
output_indices: tuple[int, ...]
equation_indices: tuple[int, ...]
equation_ids: tuple[str, ...]
evaluator_slot: int
@dataclass(frozen=True, slots=True)
class CausalIRFlowStage:
target_slots: tuple[int, ...]
scatter_compatibility_slots: tuple[int, ...]
equation_ids: tuple[str, ...]
operations: tuple[CausalIRFlowOperation, ...]
@dataclass(frozen=True, slots=True)
class CausalIRProgram:
"""Immutable callback-free structure used as the backend cache key."""
schema_version: int
canonical_slots: tuple[CausalIRCanonicalSlot, ...]
compatibility_slots: tuple[CausalIRCompatibilitySlot, ...]
reset_compatibility_slots: tuple[int, ...]
external_effort_compatibility_slots: tuple[int, ...]
effort_stages: tuple[CausalIREffortStage, ...]
flow_stages: tuple[CausalIRFlowStage, ...]
structural_signature: str
@property
def assignment_count(self) -> int:
return len(self.canonical_slots)
@property
def effort_group_count(self) -> int:
return sum(len(stage.operations) for stage in self.effort_stages)
@property
def flow_assignment_count(self) -> int:
return sum(len(stage.target_slots) for stage in self.flow_stages)
@property
def effort_scatter_count(self) -> int:
return sum(
len(operation.scatter_compatibility_slots)
for stage in self.effort_stages
for operation in stage.operations
)
@property
def eliminated_effort_replica_count(self) -> int:
return self.effort_scatter_count - self.effort_group_count
@property
def maximum_effort_stage_width(self) -> int:
return max((len(stage.operations) for stage in self.effort_stages), default=0)
@property
def maximum_flow_stage_width(self) -> int:
return max((len(stage.target_slots) for stage in self.flow_stages), default=0)
def structural_dict(self) -> dict[str, object]:
return {
"schemaVersion": self.schema_version,
"canonicalSlots": [
{
"slot": item.slot,
"id": item.id,
"variable": item.variable,
"kind": item.kind,
}
for item in self.canonical_slots
],
"compatibilitySlots": [
{"slot": item.slot, "id": item.id, "variable": item.variable}
for item in self.compatibility_slots
],
"resetCompatibilitySlots": list(self.reset_compatibility_slots),
"externalEffortCompatibilitySlots": list(
self.external_effort_compatibility_slots
),
"effortStages": [
{
"variable": stage.variable,
"operations": [
{
"opcode": operation.opcode.value,
"resultSlot": operation.result_slot,
"anchorCompatibilitySlot": (
operation.anchor_compatibility_slot
),
"scatterCompatibilitySlots": list(
operation.scatter_compatibility_slots
),
"equationId": operation.equation_id,
}
for operation in stage.operations
],
"evaluations": [
{
"opcode": evaluation.opcode.value,
"outputIndices": list(evaluation.output_indices),
"equationIndices": list(evaluation.equation_indices),
"equationIds": list(evaluation.equation_ids),
"evaluatorSlot": evaluation.evaluator_slot,
}
for evaluation in stage.evaluations
],
}
for stage in self.effort_stages
],
"flowStages": [
{
"targetSlots": list(stage.target_slots),
"scatterCompatibilitySlots": list(
stage.scatter_compatibility_slots
),
"equationIds": list(stage.equation_ids),
"operations": [
{
"opcode": operation.opcode.value,
"outputIndices": list(operation.output_indices),
"equationIndices": list(operation.equation_indices),
"equationIds": list(operation.equation_ids),
"evaluatorSlot": operation.evaluator_slot,
}
for operation in stage.operations
],
}
for stage in self.flow_stages
],
}
def calculate_structural_signature(self) -> str:
payload = json.dumps(
self.structural_dict(),
ensure_ascii=True,
separators=(",", ":"),
sort_keys=True,
).encode("utf-8")
return sha256(payload).hexdigest()
@dataclass(frozen=True, slots=True)
class CausalIRBindings:
readers: tuple[Callable[[], float], ...]
writers: tuple[Callable[[float], None], ...]
evaluators: tuple[Callable[[], object], ...]
@dataclass(slots=True)
class CausalIRWorkspace:
structural_signature: str
canonical_values: "np.ndarray[Any, Any]"
effort_residuals: "np.ndarray[Any, Any]"
effort_written: "np.ndarray[Any, Any]"
flow_values: "np.ndarray[Any, Any]"
flow_written: "np.ndarray[Any, Any]"
transaction_values: "np.ndarray[Any, Any]"
@dataclass(frozen=True, slots=True)
class CausalIRExecutionResult:
success: bool
fallback_reason: str | None
structural_signature: str
effort_assignment_count: int
flow_assignment_count: int
completed_effort_stage_count: int
completed_flow_stage_count: int
rolled_back: bool
StageObserver = Callable[
[str, int, tuple[int, ...], tuple[float, ...]],
None,
]
@dataclass(frozen=True, slots=True)
class CausalNumericIR:
"""Bound reference IR; its normal path performs no full snapshot."""
program: CausalIRProgram
bindings: CausalIRBindings
def create_workspace(self) -> CausalIRWorkspace:
try:
import numpy as np
except ImportError as exc: # pragma: no cover
raise RuntimeError("The causal numeric reference IR requires NumPy.") from exc
return CausalIRWorkspace(
structural_signature=self.program.structural_signature,
canonical_values=np.empty(
max(len(self.program.canonical_slots), 1), dtype=np.float64
),
effort_residuals=np.empty(
max(self.program.maximum_effort_stage_width, 1), dtype=np.float64
),
effort_written=np.empty(
max(self.program.maximum_effort_stage_width, 1), dtype=np.bool_
),
flow_values=np.empty(
max(self.program.maximum_flow_stage_width, 1), dtype=np.float64
),
flow_written=np.empty(
max(self.program.maximum_flow_stage_width, 1), dtype=np.bool_
),
transaction_values=np.empty(
max(len(self.program.compatibility_slots), 1), dtype=np.float64
),
)
def execute(
self,
workspace: CausalIRWorkspace,
*,
effort_variables: tuple[str, ...] = ("p",),
transactional: bool = False,
stage_observer: StageObserver | None = None,
) -> CausalIRExecutionResult:
"""Interpret the IR; transactional snapshots are audit-only."""
program = self.program
bindings = self.bindings
signature = program.structural_signature
if workspace.structural_signature != signature:
raise ValueError("Causal IR workspace belongs to a different program.")
if len(bindings.readers) != len(program.compatibility_slots) or len(
bindings.writers
) != len(program.compatibility_slots):
raise ValueError("Causal IR compatibility binding count is inconsistent.")
if any(variable not in {"p", "x", "v"} for variable in effort_variables):
return CausalIRExecutionResult(
False, "unsupportedEffortVariable", signature, 0, 0, 0, 0, False
)
snapshot_count = 0
if transactional:
try:
for slot, reader in enumerate(bindings.readers):
workspace.transaction_values[slot] = float(reader())
snapshot_count += 1
except MemoryError:
raise
except (ArithmeticError, RuntimeError, TypeError, ValueError) as exc:
return CausalIRExecutionResult(
False,
f"slotReadFailed:{type(exc).__name__}",
signature,
0,
0,
0,
0,
False,
)
effort_count = 0
flow_count = 0
completed_effort_stages = 0
completed_flow_stages = 0
def failed(reason: str) -> CausalIRExecutionResult:
rolled_back = False
if transactional:
for slot in range(snapshot_count):
bindings.writers[slot](float(workspace.transaction_values[slot]))
rolled_back = True
return CausalIRExecutionResult(
False,
reason,
signature,
effort_count,
flow_count,
completed_effort_stages,
completed_flow_stages,
rolled_back,
)
selected_efforts = frozenset(effort_variables)
for stage_index, stage in enumerate(program.effort_stages):
if stage.variable not in selected_efforts:
continue
width = len(stage.operations)
workspace.effort_written[:width] = False
for evaluation in stage.evaluations:
try:
evaluated = bindings.evaluators[evaluation.evaluator_slot]()
if evaluation.opcode is CausalIROpcode.EFFORT_DIRECT_RESIDUAL:
output = evaluation.output_indices[0]
workspace.effort_residuals[output] = float(evaluated)
workspace.effort_written[output] = True
continue
if not hasattr(evaluated, "__len__"):
raise TypeError("component evaluator returned no sequence")
for output, equation in zip(
evaluation.output_indices, evaluation.equation_indices
):
if equation >= len(evaluated):
raise IndexError("component equation disappeared")
workspace.effort_residuals[output] = float(evaluated[equation])
workspace.effort_written[output] = True
except MemoryError:
raise
except (
ArithmeticError,
IndexError,
RuntimeError,
TypeError,
ValueError,
) as exc:
return failed(f"effortEvaluationFailed:{type(exc).__name__}")
if any(not bool(workspace.effort_written[index]) for index in range(width)):
return failed("effortEvaluationCoverageMismatch")
for output, operation in enumerate(stage.operations):
try:
anchor = float(bindings.readers[operation.anchor_compatibility_slot]())
target = anchor - float(workspace.effort_residuals[output])
except MemoryError:
raise
except (
ArithmeticError,
IndexError,
RuntimeError,
TypeError,
ValueError,
) as exc:
return failed(f"effortAssignmentFailed:{type(exc).__name__}")
if not isfinite(target) or (
stage.variable == "p" and target <= PRESSURE_LOWER_BOUND_PA
):
return failed("nonFiniteOrInvalidEffortAssignment")
workspace.canonical_values[operation.result_slot] = target
for slot in operation.scatter_compatibility_slots:
bindings.writers[slot](target)
effort_count += 1
completed_effort_stages += 1
if stage_observer is not None:
try:
stage_observer(
f"effort:{stage.variable}",
stage_index,
tuple(item.result_slot for item in stage.operations),
tuple(
float(workspace.canonical_values[item.result_slot])
for item in stage.operations
),
)
except MemoryError:
raise
except Exception as exc:
return failed(f"stageObserverFailed:{type(exc).__name__}")
try:
external_finite = all(
isfinite(float(bindings.readers[slot]()))
for slot in program.external_effort_compatibility_slots
)
except MemoryError:
raise
except (ArithmeticError, RuntimeError, TypeError, ValueError) as exc:
return failed(f"externalEffortReadFailed:{type(exc).__name__}")
if not external_finite:
return failed("nonFiniteExternalEffort")
for slot in program.reset_compatibility_slots:
bindings.writers[slot](0.0)
for stage_index, stage in enumerate(program.flow_stages):
width = len(stage.target_slots)
workspace.flow_written[:width] = False
for operation in stage.operations:
try:
evaluated = bindings.evaluators[operation.evaluator_slot]()
if operation.opcode is CausalIROpcode.FLOW_DIRECT:
output = operation.output_indices[0]
workspace.flow_values[output] = float(evaluated)
workspace.flow_written[output] = True
continue
if not hasattr(evaluated, "__len__"):
raise TypeError("component evaluator returned no sequence")
for output, equation in zip(
operation.output_indices, operation.equation_indices
):
if equation >= len(evaluated):
raise IndexError("component equation disappeared")
# Targets are zero before the stage; preserve -residual.
workspace.flow_values[output] = -float(evaluated[equation])
workspace.flow_written[output] = True
except MemoryError:
raise
except (
ArithmeticError,
IndexError,
RuntimeError,
TypeError,
ValueError,
) as exc:
return failed(f"flowEvaluationFailed:{type(exc).__name__}")
if any(not bool(workspace.flow_written[index]) for index in range(width)):
return failed("flowAssignmentCoverageMismatch")
for output, (canonical, compatibility) in enumerate(
zip(stage.target_slots, stage.scatter_compatibility_slots)
):
target = float(workspace.flow_values[output])
if not isfinite(target):
return failed("nonFiniteFlowAssignment")
workspace.canonical_values[canonical] = target
bindings.writers[compatibility](target)
flow_count += 1
completed_flow_stages += 1
if stage_observer is not None:
try:
stage_observer(
"flow",
stage_index,
stage.target_slots,
tuple(float(workspace.flow_values[i]) for i in range(width)),
)
except MemoryError:
raise
except Exception as exc:
return failed(f"stageObserverFailed:{type(exc).__name__}")
return CausalIRExecutionResult(
True,
None,
signature,
effort_count,
flow_count,
completed_effort_stages,
completed_flow_stages,
False,
)
@dataclass(frozen=True, slots=True)
class CausalIRCompilation:
ir: CausalNumericIR | None
fallback_reason: str | None
@property
def supported(self) -> bool:
return self.ir is not None and self.fallback_reason is None
def _unsupported(reason: str) -> CausalIRCompilation:
return CausalIRCompilation(ir=None, fallback_reason=reason)
def _unique_slots(items: Iterable[int]) -> tuple[int, ...]:
return tuple(dict.fromkeys(int(item) for item in items))
def _compile_effort_evaluations(
operations: tuple[CausalIREffortOperation, ...],
anchor_evaluators: tuple[Callable[[], float], ...],
component_locations: dict[
str, tuple[object, Callable[[], tuple[float, ...]], int]
],
evaluators: list[Callable[[], object]],
) -> tuple[CausalIREffortEvaluation, ...]:
grouped: dict[int, list[tuple[int, int, str]]] = {}
component_callbacks: dict[int, Callable[[], tuple[float, ...]]] = {}
direct: list[tuple[int, Callable[[], float], str]] = []
for output, (operation, anchor_evaluate) in enumerate(
zip(operations, anchor_evaluators)
):
location = component_locations.get(operation.equation_id)
if location is None:
direct.append((output, anchor_evaluate, operation.equation_id))
continue
owner, evaluate, equation = location
key = id(owner)
component_callbacks[key] = evaluate
grouped.setdefault(key, []).append((output, equation, operation.equation_id))
compiled: list[CausalIREffortEvaluation] = []
for output, evaluate, equation_id in direct:
evaluator = len(evaluators)
evaluators.append(evaluate)
compiled.append(
CausalIREffortEvaluation(
CausalIROpcode.EFFORT_DIRECT_RESIDUAL,
(output,),
(),
(equation_id,),
evaluator,
)
)
for key, entries in grouped.items():
evaluator = len(evaluators)
evaluators.append(component_callbacks[key])
compiled.append(
CausalIREffortEvaluation(
CausalIROpcode.EFFORT_COMPONENT_RESIDUAL,
tuple(item[0] for item in entries),
tuple(item[1] for item in entries),
tuple(item[2] for item in entries),
evaluator,
)
)
return tuple(compiled)
def compile_causal_numeric_ir(solver: object) -> CausalIRCompilation:
"""Lower a compile-proven global plan; unsupported plans fail closed."""
if not bool(getattr(solver, "_causal_fast_path_eligible", False)):
return _unsupported(
str(
getattr(solver, "_causal_fast_path_fallback_reason", None)
or "causalProofNotAvailable"
)
)
try:
unknowns = tuple(getattr(solver, "unknowns"))
effort_plan = getattr(solver, "_causal_effort_plan_by_variable")
flow_plan = tuple(getattr(solver, "_explicit_flow_plan"))
component_plan = tuple(getattr(solver, "_component_equation_plan"))
reset_unknowns = tuple(
getattr(solver, "_explicit_flow_unknowns_by_variables")[
frozenset(("f", "m_flow"))
]
)
external_unknowns = tuple(
getattr(solver, "_causal_external_effort_unknowns")
)
except (AttributeError, KeyError, TypeError):
return _unsupported("unsupportedCausalSolverContract")
unknown_ids = tuple(str(item.id) for item in unknowns)
if len(set(unknown_ids)) != len(unknown_ids):
return _unsupported("duplicateAlgebraicUnknown")
compatibility_slot_by_id = {
unknown_id: slot for slot, unknown_id in enumerate(unknown_ids)
}
compatibility_slots = tuple(
CausalIRCompatibilitySlot(slot, unknown_id, str(unknown.variable))
for slot, (unknown_id, unknown) in enumerate(zip(unknown_ids, unknowns))
)
readers = tuple(item.read for item in unknowns)
writers = tuple(item.write for item in unknowns)
evaluators: list[Callable[[], object]] = []
canonical_slots: list[CausalIRCanonicalSlot] = []
component_locations: dict[
str, tuple[object, Callable[[], tuple[float, ...]], int]
] = {}
try:
for plan in component_plan:
for equation, template in enumerate(plan.templates):
component_locations[str(template.id)] = (
plan.component,
plan.evaluate,
equation,
)
except (AttributeError, TypeError):
return _unsupported("unsupportedComponentEvaluationContract")
effort_stages: list[CausalIREffortStage] = []
try:
for variable in ("p", "x", "v"):
operations: list[CausalIREffortOperation] = []
anchors: list[Callable[[], float]] = []
for assignment in effort_plan[variable]:
result = len(canonical_slots)
equation_id = str(assignment.anchor.equation_id)
scatter = tuple(
compatibility_slot_by_id[item.id]
for item in assignment.members
)
if not scatter or len(set(scatter)) != len(scatter):
return _unsupported("invalidEffortScatterSlots")
canonical_slots.append(
CausalIRCanonicalSlot(
result,
f"effort:{variable}:{equation_id}",
variable,
"effort_group",
)
)
operations.append(
CausalIREffortOperation(
CausalIROpcode.EFFORT_BROADCAST,
variable,
result,
compatibility_slot_by_id[assignment.anchor.unknown.id],
scatter,
equation_id,
)
)
anchors.append(assignment.anchor.evaluate)
operation_tuple = tuple(operations)
effort_stages.append(
CausalIREffortStage(
variable,
operation_tuple,
_compile_effort_evaluations(
operation_tuple,
tuple(anchors),
component_locations,
evaluators,
),
)
)
except (AttributeError, KeyError, TypeError):
return _unsupported("unsupportedEffortPlanContract")
flow_stages: list[CausalIRFlowStage] = []
try:
for stage in flow_plan:
scatter = tuple(
compatibility_slot_by_id[item.unknown.id]
for item in stage.assignments
)
equation_ids = tuple(str(item.equation_id) for item in stage.assignments)
if len(set(scatter)) != len(scatter):
return _unsupported("duplicateFlowTargetInStage")
targets: list[int] = []
for assignment in stage.assignments:
target = len(canonical_slots)
targets.append(target)
canonical_slots.append(
CausalIRCanonicalSlot(
target,
f"flow:{assignment.unknown.id}",
str(assignment.unknown.variable),
"flow_assignment",
)
)
covered: list[int] = []
operations: list[CausalIRFlowOperation] = []
for output, evaluate in stage.direct_evaluations:
output = int(output)
evaluator = len(evaluators)
evaluators.append(evaluate)
operations.append(
CausalIRFlowOperation(
CausalIROpcode.FLOW_DIRECT,
(output,),
(),
(equation_ids[output],),
evaluator,
)
)
covered.append(output)
for evaluation in stage.component_evaluations:
evaluator = len(evaluators)
evaluators.append(evaluation.evaluate)
outputs = tuple(int(item) for item in evaluation.assignment_indices)
operations.append(
CausalIRFlowOperation(
CausalIROpcode.FLOW_COMPONENT_RESIDUAL,
outputs,
tuple(int(item) for item in evaluation.equation_indices),
tuple(str(item) for item in evaluation.equation_ids),
evaluator,
)
)
covered.extend(outputs)
if sorted(covered) != list(range(len(scatter))):
return _unsupported("flowStageEvaluationCoverageMismatch")
flow_stages.append(
CausalIRFlowStage(
tuple(targets), scatter, equation_ids, tuple(operations)
)
)
except (AttributeError, IndexError, KeyError, TypeError):
return _unsupported("unsupportedFlowPlanContract")
try:
reset_slots = _unique_slots(
compatibility_slot_by_id[item.id] for item in reset_unknowns
)
external_slots = _unique_slots(
compatibility_slot_by_id[item.id] for item in external_unknowns
)
except (AttributeError, KeyError):
return _unsupported("unknownCausalBoundarySlot")
flow_scatter = tuple(
item for stage in flow_stages for item in stage.scatter_compatibility_slots
)
if len(set(flow_scatter)) != len(flow_scatter):
return _unsupported("duplicateExplicitFlowAssignment")
if set(flow_scatter) != set(reset_slots):
return _unsupported("incompleteExplicitFlowCoverage")
program = CausalIRProgram(
CAUSAL_NUMERIC_IR_SCHEMA_VERSION,
tuple(canonical_slots),
compatibility_slots,
reset_slots,
external_slots,
tuple(effort_stages),
tuple(flow_stages),
"",
)
program = replace(
program, structural_signature=program.calculate_structural_signature()
)
return CausalIRCompilation(
CausalNumericIR(
program,
CausalIRBindings(readers, writers, tuple(evaluators)),
),
None,
)
+245 -24
View File
@@ -29,6 +29,8 @@ StateTransitionHandler = Callable[
]
_MAX_STATE_TRANSITIONS_AT_SAME_TIME = 64
_MAX_RECOVERABLE_RETRIES = 16
_RECOVERABLE_RETRY_FACTOR = 0.5
class IntegrationCancelled(Exception):
@@ -48,6 +50,29 @@ class SolveIVPConfig:
first_step: float | None = None
@dataclass(frozen=True)
class RecoverableRetryDiagnostics:
"""One recoverable trial failure and the step cap chosen for its retry."""
phase: Literal["constructor", "step", "solver-status"]
attempted_step: float
reason: str
next_max_step: float | None = None
next_first_step: float | None = None
def as_dict(self) -> dict[str, object]:
result: dict[str, object] = {
"phase": self.phase,
"attemptedStep": self.attempted_step,
"reason": self.reason,
}
if self.next_max_step is not None:
result["nextMaxStep"] = self.next_max_step
if self.next_first_step is not None:
result["nextFirstStep"] = self.next_first_step
return result
@dataclass(frozen=True)
class SolverSegmentDiagnostics:
"""Work performed by implicit solver instances inside one event segment."""
@@ -61,6 +86,7 @@ class SolverSegmentDiagnostics:
accepted_step_count: int = 0
solver_start_count: int = 0
state_transition_count: int = 0
state_transition_times: tuple[float, ...] = ()
recoverable_retry_count: int = 0
jacobian_evaluation_count: int = 0
jacobian_full_build_count: int = 0
@@ -72,9 +98,10 @@ class SolverSegmentDiagnostics:
exact_column_build_count: int = 0
exact_column_fallback_count: int = 0
jacobian_assembly_seconds: float = 0.0
recoverable_retries: tuple[RecoverableRetryDiagnostics, ...] = ()
def as_dict(self) -> dict[str, float | int]:
result: dict[str, float | int] = {
def as_dict(self) -> dict[str, object]:
result: dict[str, object] = {
"startTime": self.start_time,
"requestedStopTime": self.requested_stop_time,
"simulatedUntil": self.simulated_until,
@@ -86,6 +113,14 @@ class SolverSegmentDiagnostics:
"stateTransitionCount": self.state_transition_count,
"recoverableRetryCount": self.recoverable_retry_count,
}
if self.state_transition_times:
result["stateTransitionTimes"] = list(
self.state_transition_times
)
if self.recoverable_retries:
result["recoverableRetries"] = [
retry.as_dict() for retry in self.recoverable_retries
]
if (
self.jacobian_evaluation_count
or self.finite_difference_rhs_evaluation_count
@@ -143,6 +178,80 @@ def _jacobian_diagnostic_snapshot(
}
def _positive_finite_step(value: object) -> float | None:
if value is None:
return None
try:
candidate = abs(float(value))
except (TypeError, ValueError, OverflowError):
return None
return candidate if candidate > 0.0 and math.isfinite(candidate) else None
def _smallest_positive_finite_step(*values: object) -> float:
"""Return a conservative step bound from configuration candidates."""
candidates = [
candidate
for value in values
if (candidate := _positive_finite_step(value)) is not None
]
if not candidates:
raise ValueError("No positive finite integration step is available.")
return min(candidates)
def _solver_attempted_step(
solver: object,
*,
segment_max_step: float,
remaining_interval: float,
) -> float:
"""Snapshot the real trial scale before calling ``solver.step()``.
SciPy exposes the proposed step as ``h_abs``. ``step_size`` is the prior
accepted step, so it is only a fallback for solvers without a valid
``h_abs``; it must not reduce an otherwise valid failed-trial estimate.
"""
configured_cap = _smallest_positive_finite_step(
segment_max_step,
remaining_interval,
)
for attribute in ("h_abs", "step_size"):
try:
candidate = _positive_finite_step(
getattr(solver, attribute, None)
)
except Exception:
# A third-party OdeSolver may implement these as fragile
# properties. The configured cap remains a safe fallback.
continue
if candidate is not None:
return min(candidate, configured_cap)
return configured_cap
def _recoverable_retry_steps(
attempted_step: float,
*,
last_accepted_time: float,
) -> tuple[float, float] | None:
"""Return strictly smaller max/first steps, or None at machine precision."""
next_step = _RECOVERABLE_RETRY_FACTOR * attempted_step
minimum_step = 64.0 * math.ulp(max(abs(last_accepted_time), 1.0))
if (
not math.isfinite(next_step)
or next_step <= minimum_step
or next_step >= attempted_step
):
return None
# This first step is intentionally one-shot. Keeping it equal to the new
# cap makes both controls strictly smaller than the failed trial scale.
return next_step, next_step
@dataclass(frozen=True)
class ODESolution:
t: list[float]
@@ -751,13 +860,18 @@ def _integrate_scipy_stepwise(
segment_max_step = float(config.max_step)
recoverable_retry_count = 0
last_recoverable_error: RecoverableTrialStateError | None = None
retry_first_step: float | None = None
segment_nfev = 0
segment_njev = 0
segment_nlu = 0
segment_accepted_steps = 0
segment_solver_starts = 0
segment_state_transitions = 0
segment_state_transition_times: list[float] = []
segment_recoverable_retries = 0
segment_recoverable_retry_diagnostics: list[
RecoverableRetryDiagnostics
] = []
jacobian_work_start = _jacobian_diagnostic_snapshot(implicit_jac)
while has_integration_interval and last_accepted_time < integration_end:
@@ -777,8 +891,8 @@ def _integrate_scipy_stepwise(
elif jac_sparsity is not None:
solver_options["jac_sparsity"] = jac_sparsity
requested_first_step = (
0.1 * segment_max_step
if last_recoverable_error is not None
retry_first_step
if retry_first_step is not None
else config.first_step
)
if requested_first_step is not None:
@@ -786,8 +900,12 @@ def _integrate_scipy_stepwise(
requested_first_step,
integration_end - last_accepted_time,
)
try:
constructor_attempted_step = _smallest_positive_finite_step(
segment_max_step,
integration_end - last_accepted_time,
solver_options.get("first_step"),
)
start_segment = getattr(implicit_jac, "start_segment", None)
if start_segment is not None:
start_segment()
@@ -806,14 +924,33 @@ def _integrate_scipy_stepwise(
recoverable_retry_count += 1
segment_recoverable_retries += 1
last_recoverable_error = exc
next_step = 0.5 * segment_max_step
minimum_step = 64.0 * math.ulp(max(abs(last_accepted_time), 1.0))
if recoverable_retry_count > 16 or next_step <= minimum_step:
retry_steps = (
_recoverable_retry_steps(
constructor_attempted_step,
last_accepted_time=last_accepted_time,
)
if recoverable_retry_count <= _MAX_RECOVERABLE_RETRIES
else None
)
segment_recoverable_retry_diagnostics.append(
RecoverableRetryDiagnostics(
phase="constructor",
attempted_step=constructor_attempted_step,
reason=str(exc),
next_max_step=(
retry_steps[0] if retry_steps is not None else None
),
next_first_step=(
retry_steps[1] if retry_steps is not None else None
),
)
)
if retry_steps is None:
status = "failed"
message = str(exc)
error = exc
break
segment_max_step = next_step
segment_max_step, retry_first_step = retry_steps
continue
except Exception as exc:
status = "failed"
@@ -836,6 +973,13 @@ def _integrate_scipy_stepwise(
step_start_time = last_accepted_time
step_start_state = list(last_accepted_state)
try:
attempted_step = _solver_attempted_step(
solver,
segment_max_step=segment_max_step,
remaining_interval=(
integration_end - last_accepted_time
),
)
step_message = solver.step()
except IntegrationCancelled:
status = "cancelled"
@@ -847,17 +991,38 @@ def _integrate_scipy_stepwise(
recoverable_retry_count += 1
segment_recoverable_retries += 1
last_recoverable_error = exc
attempted_step = segment_max_step
next_step = 0.5 * attempted_step
minimum_step = 64.0 * math.ulp(
max(abs(last_accepted_time), 1.0)
retry_steps = (
_recoverable_retry_steps(
attempted_step,
last_accepted_time=last_accepted_time,
)
if recoverable_retry_count
<= _MAX_RECOVERABLE_RETRIES
else None
)
if recoverable_retry_count > 16 or next_step <= minimum_step:
segment_recoverable_retry_diagnostics.append(
RecoverableRetryDiagnostics(
phase="step",
attempted_step=attempted_step,
reason=str(exc),
next_max_step=(
retry_steps[0]
if retry_steps is not None
else None
),
next_first_step=(
retry_steps[1]
if retry_steps is not None
else None
),
)
)
if retry_steps is None:
status = "failed"
message = str(exc)
error = exc
break
segment_max_step = next_step
segment_max_step, retry_first_step = retry_steps
restart_after_recoverable = True
break
except Exception as exc:
@@ -871,21 +1036,62 @@ def _integrate_scipy_stepwise(
if last_recoverable_error is not None:
recoverable_retry_count += 1
segment_recoverable_retries += 1
next_step = 0.5 * segment_max_step
minimum_step = 64.0 * math.ulp(
max(abs(last_accepted_time), 1.0)
retry_steps = (
_recoverable_retry_steps(
attempted_step,
last_accepted_time=last_accepted_time,
)
if recoverable_retry_count
<= _MAX_RECOVERABLE_RETRIES
else None
)
if (
recoverable_retry_count <= 16
and next_step > minimum_step
):
segment_max_step = next_step
failure_reason = str(
step_message or last_recoverable_error
)
segment_recoverable_retry_diagnostics.append(
RecoverableRetryDiagnostics(
phase="solver-status",
attempted_step=attempted_step,
reason=failure_reason,
next_max_step=(
retry_steps[0]
if retry_steps is not None
else None
),
next_first_step=(
retry_steps[1]
if retry_steps is not None
else None
),
)
)
if retry_steps is not None:
segment_max_step, retry_first_step = retry_steps
restart_after_recoverable = True
break
status = "failed"
message = str(step_message or "Integration step failed.")
break
# A returned running/finished status means this step was
# accepted. Any prior recoverable failure is now historical:
# it must not influence an event restart or an ordinary later
# solver failure. The reduced cap is local to the failed
# trial: after one accepted retry step, let this solver grow
# adaptively again and ensure a later event restart receives
# the configured maximum. The retry-specific first step is
# likewise strictly one-shot.
if retry_first_step is not None:
segment_max_step = float(config.max_step)
try:
solver.max_step = segment_max_step
except (AttributeError, TypeError, ValueError):
# Third-party OdeSolver-compatible test doubles may not
# expose a writable cap. SciPy's supported solvers do.
pass
last_recoverable_error = None
retry_first_step = None
recoverable_retry_count = 0
segment_accepted_steps += 1
step_end_time = float(solver.t)
step_end_state = [float(value) for value in solver.y]
@@ -940,6 +1146,9 @@ def _integrate_scipy_stepwise(
if transition is not None:
segment_state_transitions += 1
segment_state_transition_times.append(
float(transition.time)
)
try:
same_time_transition_count = (
_next_same_time_transition_count(
@@ -997,7 +1206,6 @@ def _integrate_scipy_stepwise(
last_accepted_time = step_end_time
last_accepted_state = step_end_state
recoverable_retry_count = 0
reported_time = (
float(segment_end)
if is_breakpoint and solver.status == "finished"
@@ -1057,7 +1265,13 @@ def _integrate_scipy_stepwise(
accepted_step_count=segment_accepted_steps,
solver_start_count=segment_solver_starts,
state_transition_count=segment_state_transitions,
state_transition_times=tuple(
segment_state_transition_times
),
recoverable_retry_count=segment_recoverable_retries,
recoverable_retries=tuple(
segment_recoverable_retry_diagnostics
),
jacobian_evaluation_count=int(
jacobian_work["jacobianEvaluationCount"]
),
@@ -1150,6 +1364,7 @@ def integrate_ode(
state_transition_handler: StateTransitionHandler | None = None,
jac_sparsity=None,
jac: JacobianCallable | None = None,
recoverable_trial_retries: bool = False,
):
"""Integrate an ODE, optionally restarting at equation discontinuities.
@@ -1161,6 +1376,11 @@ def integrate_ode(
interpolant. When it returns a transition, samples before the event retain
the pre-event trajectory, the reset state is stored at the event, and a fresh
solver continues from that state.
``recoverable_trial_retries`` opts an eventless/cancellation-free caller
into the stepwise path so a ``RecoverableTrialStateError`` can rebuild the
solver from its last accepted state. It defaults to false to preserve the
direct ``solve_ivp`` path for ordinary callers.
"""
if (
@@ -1209,6 +1429,7 @@ def integrate_ode(
cancel_check is not None
or normalized_breakpoints
or state_transition_handler is not None
or recoverable_trial_retries
):
return _integrate_scipy_stepwise(
rhs,
+20
View File
@@ -60,6 +60,16 @@ class StreamResolver:
for component in self._components
if not isinstance(component, DynamicComponent)
)
# State ownership and pressure-flow stream sensitivity are independent
# classifications. Compile this hook by behavior so algebraic
# components such as PNL00R receive their upstream-temperature
# references without dispatching a no-op to every component at runtime.
self._flow_temperature_reference_components = tuple(
component
for component in self._components
if type(component).update_flow_temperature_references
is not Component.update_flow_temperature_references
)
self._ports = tuple(
(component.name, port_name, port)
for component in self._components
@@ -122,6 +132,16 @@ class StreamResolver:
)
return values
@profile_phase("simulation.refresh", minimum_mode="audit")
def refresh_flow_temperature_references(self) -> None:
"""Refresh pressure-flow property inputs without changing stream outflows."""
connected = self.connected_temperature_reference_enthalpies()
for component in self._flow_temperature_reference_components:
component.update_flow_temperature_references(
connected[component.name]
)
@profile_phase("simulation.refresh", minimum_mode="audit")
def _refresh_dynamic_components(self) -> None:
for component in self._dynamic_components:
+168 -11
View File
@@ -1,11 +1,13 @@
"""Proof-gated tangent columns for the three-piston reference network.
"""Proof-gated tangent columns for supported piston branch networks.
This module is deliberately narrower than the generic algebraic solver. It
only compiles a tangent provider after proving the state layout, component
types, physical connections, and causal execution plan used by the committed
three-piston XML. A failed proof leaves the ordinary seed-0 numerical
Jacobian in control; a runtime mode boundary requests the same one-build
fallback through :class:`ExactColumnsUnavailable`.
types, physical connections, and causal execution plan used by every selected
piston branch. The legacy three-piston entry point remains available for its
committed fixture, while the topology-driven entry point discovers any number
of branches without depending on component names. A failed proof leaves the
ordinary seed-0 numerical Jacobian in control; a runtime mode boundary requests
the same one-build fallback through :class:`ExactColumnsUnavailable`.
"""
from __future__ import annotations
@@ -57,6 +59,7 @@ class ThreePistonBranch:
chamber: object
pipe: object
contact: object
chamber_connection_port: str
velocity_index: int
position_index: int
@@ -91,7 +94,7 @@ def _failed(reason: str) -> ThreePistonTangentCompilation:
class ThreePistonTangentProvider:
"""Batched six-direction provider compiled for one system instance."""
"""Batched selected-branch provider compiled for one system instance."""
def __init__(
self,
@@ -611,10 +614,17 @@ class ThreePistonTangentProvider:
return out
def compile_three_piston_tangent_provider(
@dataclass(frozen=True)
class _PistonBranchSpec:
names: tuple[str, str, str, str, str]
chamber_connection_port: str
def _compile_named_piston_tangent_provider(
system: "GenericFluidSystem",
branch_specs: Sequence[_PistonBranchSpec],
) -> ThreePistonTangentCompilation:
"""Compile the proof-gated target provider, or return a stable reason."""
"""Compile a named, topology-proven set of supported piston branches."""
solver = system.pressure_flow_solver
if not solver.causal_fast_path_eligible:
@@ -651,7 +661,8 @@ def compile_three_piston_tangent_provider(
"amesim_lstp00a",
)
branches: list[ThreePistonBranch] = []
for names in _TARGET_BRANCH_NAMES:
for branch_spec in branch_specs:
names = branch_spec.names
try:
components = tuple(system.network.components[name] for name in names)
except KeyError:
@@ -671,6 +682,9 @@ def compile_three_piston_tangent_provider(
chamber=chamber,
pipe=pipe,
contact=contact,
chamber_connection_port=(
branch_spec.chamber_connection_port
),
velocity_index=offset,
position_index=offset + 1,
)
@@ -681,7 +695,15 @@ def compile_three_piston_tangent_provider(
required_pairs.update(
{
frozenset((Endpoint(branch.mass.name, "port_1"), Endpoint(branch.piston.name, "port_2"))),
frozenset((Endpoint(branch.piston.name, "port_1"), Endpoint(branch.chamber.name, "port_3"))),
frozenset(
(
Endpoint(branch.piston.name, "port_1"),
Endpoint(
branch.chamber.name,
branch.chamber_connection_port,
),
)
),
frozenset((Endpoint(branch.chamber.name, "port_1"), Endpoint(branch.pipe.name, "port_1"))),
frozenset((Endpoint(branch.piston.name, "port_5"), Endpoint(branch.contact.name, "port_1"))),
}
@@ -790,7 +812,7 @@ def compile_three_piston_tangent_provider(
# Secondary pressure blocks may contain the same causal flow coordinates,
# so membership alone is not evidence of a stream derivative. The direct
# enthalpy reach proof below, plus the runtime dynamic-owner gate, is the
# relevant condition for this target-specific program.
# relevant condition for this proof-gated branch program.
neighbor_by_endpoint: dict[Endpoint, Endpoint] = {}
for connection in system.network.connections:
if connection.kind != "physical":
@@ -848,3 +870,138 @@ def compile_three_piston_tangent_provider(
provider,
reached_assignment_count=len(reached_assignments),
)
def _physical_neighbor_map(
system: "GenericFluidSystem",
) -> dict[Endpoint, Endpoint]:
"""Return the one-to-one physical connector map proved by the network."""
neighbors: dict[Endpoint, Endpoint] = {}
for connection in system.network.connections:
if connection.kind != "physical":
continue
first, second = connection.endpoints
# SimulationNetwork already rejects multiply connected physical ports.
# Retain a defensive gate because this compiler may also be called by
# custom network builders in tests or downstream applications.
if first in neighbors or second in neighbors:
raise ValueError("Physical endpoint has more than one connection.")
neighbors[first] = second
neighbors[second] = first
return neighbors
def _discover_supported_piston_branch_specs(
system: "GenericFluidSystem",
) -> tuple[_PistonBranchSpec, ...] | ThreePistonTangentCompilation:
"""Discover every complete catalog piston branch by type and port topology.
A PNRP17 is the unambiguous root: its mechanical piston-side port must be
driven by a singleton MECMAS21 coordinate, its pneumatic port must feed a
PNCH012 whose first port feeds PNL0001, and its rod-side port must meet an
LSTP00A contact. If even one PNRP17 is only partially supported, reject the
batch with a stable reason instead of silently omitting derivative columns.
"""
try:
neighbors = _physical_neighbor_map(system)
except ValueError:
return _failed("unsupportedPistonBranchTopology:multipleConnection")
components = system.network.components
def model_type(endpoint: Endpoint | None) -> str | None:
if endpoint is None:
return None
return getattr(components[endpoint.component], "MODEL_TYPE", None)
pistons = tuple(
component
for component in components.values()
if getattr(component, "MODEL_TYPE", None) == "amesim_pnrp17"
)
if not pistons:
return _failed("supportedPistonBranchMissing")
specs: list[_PistonBranchSpec] = []
for piston in pistons:
mass_endpoint = neighbors.get(Endpoint(piston.name, "port_2"))
if (
model_type(mass_endpoint) != "amesim_mecmas21"
or mass_endpoint is None
or mass_endpoint.port != "port_1"
):
return _failed("unsupportedPistonBranchTopology:mass")
chamber_endpoint = neighbors.get(Endpoint(piston.name, "port_1"))
if model_type(chamber_endpoint) != "amesim_pnch012":
return _failed("unsupportedPistonBranchTopology:chamber")
assert chamber_endpoint is not None
pipe_endpoint = neighbors.get(
Endpoint(chamber_endpoint.component, "port_1")
)
if (
model_type(pipe_endpoint) != "amesim_pnl0001"
or pipe_endpoint is None
or pipe_endpoint.port != "port_1"
):
return _failed("unsupportedPistonBranchTopology:pipe")
contact_endpoint = neighbors.get(Endpoint(piston.name, "port_5"))
if (
model_type(contact_endpoint) != "amesim_lstp00a"
or contact_endpoint is None
or contact_endpoint.port != "port_1"
):
return _failed("unsupportedPistonBranchTopology:contact")
specs.append(
_PistonBranchSpec(
names=(
mass_endpoint.component,
piston.name,
chamber_endpoint.component,
pipe_endpoint.component,
contact_endpoint.component,
),
chamber_connection_port=chamber_endpoint.port,
)
)
# Port uniqueness already proves unique pistons and masses, but explicitly
# reject a custom multi-port chamber/contact/pipe shared by two roots. The
# tangent propagation assumes one geometry seed per selected state owner.
for role_index in range(5):
if len({spec.names[role_index] for spec in specs}) != len(specs):
return _failed("unsupportedPistonBranchTopology:sharedComponent")
return tuple(specs)
def compile_supported_piston_tangent_provider(
system: "GenericFluidSystem",
) -> ThreePistonTangentCompilation:
"""Compile all name-independent, topology-supported piston branches."""
discovered = _discover_supported_piston_branch_specs(system)
if isinstance(discovered, ThreePistonTangentCompilation):
return discovered
return _compile_named_piston_tangent_provider(system, discovered)
def compile_three_piston_tangent_provider(
system: "GenericFluidSystem",
) -> ThreePistonTangentCompilation:
"""Compile the committed legacy three-piston target by its stable names."""
return _compile_named_piston_tangent_provider(
system,
tuple(
_PistonBranchSpec(
names=names,
chamber_connection_port="port_3",
)
for names in _TARGET_BRANCH_NAMES
),
)
+567
View File
@@ -0,0 +1,567 @@
from __future__ import annotations
from collections.abc import Callable, Sequence
from copy import copy
from dataclasses import dataclass, replace
from app.simulation.core.errors import RecoverableTrialStateError
from app.simulation.core.ports import PortState
_STREAM_CACHE_ATTRIBUTE_NAMES = frozenset(
{
"_connected_h",
"temperature_reference_h",
}
)
def _is_stream_cache_attribute(name: str) -> bool:
"""Return whether an attribute belongs to the stream/temperature replay state.
Catalog components currently use ``_connected_h`` and
``temperature_reference_h``. The name-based extension keeps conservative
third-party caches recoverable without copying an entire component graph.
Components with opaque cache names can provide the explicit hooks documented
by :class:`ThermofluidTransactionPlan`.
"""
lowered = name.lower()
return (
name in _STREAM_CACHE_ATTRIBUTE_NAMES
or lowered.startswith("_stream_")
or "connected_h" in lowered
or "connected_enthalpy" in lowered
or "temperature_reference" in lowered
)
def _copy_cache_value(value: object) -> object:
"""Shallow-copy a stream cache without traversing the component graph."""
if isinstance(value, (dict, list, set, bytearray)):
return copy(value)
return value
@dataclass(frozen=True)
class ThermofluidWorstPort:
component: str
port: str
value: float
signed_delta: float
def as_dict(self) -> dict[str, object]:
return {
"component": self.component,
"port": self.port,
"value": self.value,
"signedDelta": self.signed_delta,
}
@dataclass(frozen=True)
class ThermofluidIterationDelta:
iteration: int
max_delta: float
scale: float
tolerance: float
worst_port: ThermofluidWorstPort | None
def as_dict(self) -> dict[str, object]:
return {
"iteration": self.iteration,
"maxDelta": self.max_delta,
"scale": self.scale,
"tolerance": self.tolerance,
"worstPort": (
self.worst_port.as_dict()
if self.worst_port is not None
else None
),
}
@dataclass(frozen=True)
class ThermofluidClosureSuccess:
rhs_time: float
iterations: int
max_delta: float
scale: float
tolerance: float
worst_port: ThermofluidWorstPort | None
@classmethod
def from_iteration(
cls,
rhs_time: float,
delta: ThermofluidIterationDelta,
) -> ThermofluidClosureSuccess:
return cls(
rhs_time=float(rhs_time),
iterations=delta.iteration,
max_delta=delta.max_delta,
scale=delta.scale,
tolerance=delta.tolerance,
worst_port=delta.worst_port,
)
def as_dict(self) -> dict[str, object]:
return {
"rhsTime": self.rhs_time,
"iterations": self.iterations,
"maxDelta": self.max_delta,
"scale": self.scale,
"tolerance": self.tolerance,
"worstPort": (
self.worst_port.as_dict()
if self.worst_port is not None
else None
),
}
@dataclass(frozen=True)
class ThermofluidClosureFailure:
failed_rhs_time: float
iterations: int
delta_tail: tuple[ThermofluidIterationDelta, ...]
max_delta: float
scale: float
tolerance: float
worst_port: ThermofluidWorstPort | None
failure_count: int = 0
@classmethod
def from_iterations(
cls,
failed_rhs_time: float,
deltas: Sequence[ThermofluidIterationDelta],
*,
tail_limit: int = 8,
) -> ThermofluidClosureFailure:
if not deltas:
raise ValueError("A thermofluid failure requires iteration diagnostics.")
final = deltas[-1]
return cls(
failed_rhs_time=float(failed_rhs_time),
iterations=final.iteration,
delta_tail=tuple(deltas[-tail_limit:]),
max_delta=final.max_delta,
scale=final.scale,
tolerance=final.tolerance,
worst_port=final.worst_port,
)
def as_dict(self) -> dict[str, object]:
return {
"failedRhsTime": self.failed_rhs_time,
"iterations": self.iterations,
"deltaTail": [item.as_dict() for item in self.delta_tail],
"maxDelta": self.max_delta,
"scale": self.scale,
"tolerance": self.tolerance,
"worstPort": (
self.worst_port.as_dict()
if self.worst_port is not None
else None
),
"failureCount": self.failure_count,
}
class ThermofluidClosureError(RecoverableTrialStateError):
"""Recoverable exhaustion of the stream/pressure-flow fixed point.
Stream propagation failures and algebraic-solver failures intentionally
retain their original exception types: rollback is still applied, but a
smaller ODE step is not known to repair those structural/numerical errors.
"""
def __init__(self, diagnostics: ThermofluidClosureFailure) -> None:
super().__init__(
"Stream enthalpy and pressure-flow coupling did not converge "
f"after {diagnostics.iterations} iterations at "
f"t={diagnostics.failed_rhs_time:.17g}."
)
self.diagnostics = diagnostics
class ThermofluidClosureDiagnostics:
"""Run-level RHS outcomes; maintenance/postprocessing calls do not write it."""
def __init__(self) -> None:
self.failure_count = 0
self.last_failure: ThermofluidClosureFailure | None = None
self.last_success: ThermofluidClosureSuccess | None = None
def record_success(self, success: ThermofluidClosureSuccess) -> None:
self.last_success = success
def record_failure(
self,
failure: ThermofluidClosureFailure,
) -> ThermofluidClosureFailure:
self.failure_count += 1
recorded = replace(failure, failure_count=self.failure_count)
self.last_failure = recorded
return recorded
def as_dict(self) -> dict[str, object]:
return {
"failureCount": self.failure_count,
"lastFailure": (
self.last_failure.as_dict()
if self.last_failure is not None
else None
),
"lastSuccess": (
self.last_success.as_dict()
if self.last_success is not None
else None
),
}
@dataclass(frozen=True)
class _PortValueBinding:
component_name: str
port_name: str
state: PortState
variable: str
@dataclass(frozen=True)
class _PortFieldPlan:
variable: str
states: tuple[PortState, ...]
@dataclass(frozen=True)
class _FlowBinding:
component_name: str
port_name: str
state: PortState
@dataclass(frozen=True)
class _ComponentCacheBinding:
component: object
attribute_names: tuple[str, ...]
attribute_name_set: frozenset[str]
snapshot_hook: Callable[[], object] | None
restore_hook: Callable[[object], None] | None
@dataclass
class ThermofluidTransactionSnapshot:
plan: ThermofluidTransactionPlan
port_values: tuple[list[float], ...]
component_cache_values: tuple[list[object], ...]
custom_cache_values: list[object | None]
diagnostic_values: list[object]
def restore(self) -> None:
plan = self.plan
plan._restore_port_values(self.port_values)
for binding, values, custom_value in zip(
plan.component_cache_bindings,
self.component_cache_values,
self.custom_cache_values,
):
component = binding.component
for name in tuple(getattr(component, "__dict__", {})):
if (
name.startswith("_causal_")
or _is_stream_cache_attribute(name)
) and name not in binding.attribute_name_set:
delattr(component, name)
for name, value in zip(binding.attribute_names, values):
setattr(component, name, _copy_cache_value(value))
if binding.restore_hook is not None:
binding.restore_hook(custom_value)
for owner, value in zip(
plan.diagnostic_owners,
self.diagnostic_values,
):
owner.last_diagnostics = value
class ThermofluidTransactionPlan:
"""Compiled, lightweight rollback boundary for one Generic RHS closure.
It snapshots active physical-port values, catalog stream-temperature caches,
component ``_causal_*`` seed fields, and resolver/solver last diagnostics.
A custom stream-aware component with an opaque mutable cache can implement
both ``snapshot_thermofluid_closure_cache()`` and
``restore_thermofluid_closure_cache(snapshot)``; these hooks are invoked in
addition to the standard name-based cache capture.
"""
def __init__(
self,
*,
port_value_bindings: tuple[_PortValueBinding, ...],
port_field_plans: tuple[_PortFieldPlan, ...],
flow_bindings: tuple[_FlowBinding, ...],
component_cache_bindings: tuple[_ComponentCacheBinding, ...],
component_count: int,
diagnostic_owners: tuple[object, ...],
) -> None:
self.port_value_bindings = port_value_bindings
self.port_field_plans = port_field_plans
self.flow_bindings = flow_bindings
self.component_cache_bindings = component_cache_bindings
self.component_count = component_count
self.diagnostic_owners = diagnostic_owners
self._snapshot = ThermofluidTransactionSnapshot(
plan=self,
port_values=tuple(
[0.0] * len(field.states)
for field in port_field_plans
),
component_cache_values=tuple(
[None] * len(binding.attribute_names)
for binding in component_cache_bindings
),
custom_cache_values=[None] * len(component_cache_bindings),
diagnostic_values=[None] * len(diagnostic_owners),
)
@classmethod
def compile(
cls,
network: object,
*,
diagnostic_owners: Sequence[object] = (),
) -> ThermofluidTransactionPlan:
components = tuple(getattr(network, "components").values())
port_value_bindings: list[_PortValueBinding] = []
port_states_by_variable: dict[str, list[PortState]] = {}
flow_bindings: list[_FlowBinding] = []
component_cache_bindings: list[_ComponentCacheBinding] = []
for component in components:
active_definitions = tuple(
definition
for definition in component.active_port_definitions
if definition.kind == "physical"
)
for definition in active_definitions:
state = component.get_port(definition.name)
flow_bindings.append(
_FlowBinding(component.name, definition.name, state)
)
for variable in definition.variables:
port_states_by_variable.setdefault(variable.name, []).append(state)
port_value_bindings.append(
_PortValueBinding(
component.name,
definition.name,
state,
variable.name,
)
)
attribute_names = tuple(
name
for name in getattr(component, "__dict__", {})
if name.startswith("_causal_")
or _is_stream_cache_attribute(name)
)
snapshot_hook = getattr(
component,
"snapshot_thermofluid_closure_cache",
None,
)
restore_hook = getattr(
component,
"restore_thermofluid_closure_cache",
None,
)
hooks_are_available = callable(snapshot_hook) and callable(restore_hook)
if attribute_names or hooks_are_available:
component_cache_bindings.append(
_ComponentCacheBinding(
component=component,
attribute_names=attribute_names,
attribute_name_set=frozenset(attribute_names),
snapshot_hook=(snapshot_hook if hooks_are_available else None),
restore_hook=(restore_hook if hooks_are_available else None),
)
)
owners = tuple(
dict.fromkeys(
owner
for owner in diagnostic_owners
if hasattr(owner, "last_diagnostics")
)
)
return cls(
port_value_bindings=tuple(port_value_bindings),
port_field_plans=tuple(
_PortFieldPlan(variable, tuple(states))
for variable, states in port_states_by_variable.items()
),
flow_bindings=tuple(flow_bindings),
component_cache_bindings=tuple(component_cache_bindings),
component_count=len(components),
diagnostic_owners=owners,
)
def capture(self) -> ThermofluidTransactionSnapshot:
# GenericFluidSystem executes one RHS serially. Reuse one compiled
# workspace rather than allocating a snapshot object and several outer
# tuples at every successful trial point.
snapshot = self._snapshot
self._capture_port_values(snapshot.port_values)
for binding, values in zip(
self.component_cache_bindings,
snapshot.component_cache_values,
):
for position, name in enumerate(binding.attribute_names):
values[position] = _copy_cache_value(
getattr(binding.component, name)
)
for position, binding in enumerate(self.component_cache_bindings):
snapshot.custom_cache_values[position] = (
binding.snapshot_hook()
if binding.snapshot_hook is not None
else None
)
for position, owner in enumerate(self.diagnostic_owners):
snapshot.diagnostic_values[position] = owner.last_diagnostics
return snapshot
def _capture_port_values(
self,
workspaces: tuple[list[float], ...],
) -> None:
for field, values in zip(self.port_field_plans, workspaces):
variable = field.variable
states = field.states
if variable == "p":
for position, state in enumerate(states):
values[position] = state.p
elif variable == "m_flow":
for position, state in enumerate(states):
values[position] = state.m_flow
elif variable == "h_outflow":
for position, state in enumerate(states):
values[position] = state.h_outflow
elif variable == "volume":
for position, state in enumerate(states):
values[position] = state.volume
elif variable == "volume_flow":
for position, state in enumerate(states):
values[position] = state.volume_flow
elif variable == "x":
for position, state in enumerate(states):
values[position] = state.x
elif variable == "v":
for position, state in enumerate(states):
values[position] = state.v
elif variable == "f":
for position, state in enumerate(states):
values[position] = state.f
else:
for position, state in enumerate(states):
values[position] = getattr(state, variable)
def _restore_port_values(
self,
workspaces: tuple[list[float], ...],
) -> None:
for field, values in zip(self.port_field_plans, workspaces):
variable = field.variable
states = field.states
if variable == "p":
for state, value in zip(states, values):
state.p = value
elif variable == "m_flow":
for state, value in zip(states, values):
state.m_flow = value
elif variable == "h_outflow":
for state, value in zip(states, values):
state.h_outflow = value
elif variable == "volume":
for state, value in zip(states, values):
state.volume = value
elif variable == "volume_flow":
for state, value in zip(states, values):
state.volume_flow = value
elif variable == "x":
for state, value in zip(states, values):
state.x = value
elif variable == "v":
for state, value in zip(states, values):
state.v = value
elif variable == "f":
for state, value in zip(states, values):
state.f = value
else:
for state, value in zip(states, values):
setattr(state, variable, value)
def flow_values(self) -> tuple[float, ...]:
return tuple(float(binding.state.m_flow) for binding in self.flow_bindings)
def measure_flow_delta(
self,
previous: Sequence[float],
*,
iteration: int,
relative_tolerance: float,
) -> ThermofluidIterationDelta:
current = self.flow_values()
scale = max(
(abs(value) for value in (*previous, *current)),
default=1.0,
)
scale = max(scale, 1.0)
worst_index = -1
worst_signed_delta = 0.0
max_delta = 0.0
for index, (old, new) in enumerate(zip(previous, current)):
signed_delta = new - old
magnitude = abs(signed_delta)
if magnitude > max_delta:
worst_index = index
worst_signed_delta = signed_delta
max_delta = magnitude
worst_port = None
if worst_index >= 0:
binding = self.flow_bindings[worst_index]
worst_port = ThermofluidWorstPort(
component=binding.component_name,
port=binding.port_name,
value=current[worst_index],
signed_delta=worst_signed_delta,
)
return ThermofluidIterationDelta(
iteration=int(iteration),
max_delta=max_delta,
scale=scale,
tolerance=float(relative_tolerance) * scale,
worst_port=worst_port,
)
def diagnostics(self) -> dict[str, int]:
stream_cache_slot_count = sum(
len(binding.attribute_names)
for binding in self.component_cache_bindings
)
return {
"physicalPortValueSlotCount": len(self.port_value_bindings),
"physicalFlowPortCount": len(self.flow_bindings),
"componentCount": self.component_count,
"cacheBindingCount": len(self.component_cache_bindings),
"streamAndCausalCacheSlotCount": stream_cache_slot_count,
"customCacheHookCount": sum(
binding.snapshot_hook is not None
for binding in self.component_cache_bindings
),
"diagnosticOwnerCount": len(self.diagnostic_owners),
}