from __future__ import annotations from collections.abc import Callable, Mapping from dataclasses import dataclass, replace from math import isfinite from app.simulation.core.equations import EquationResidual from app.simulation.solvers.algebraic import ( PRESSURE_LOWER_BOUND_PA, AlgebraicSolveDiagnostics, AlgebraicUnknown, ExplicitFlowStage, PressureFlowSolver, ) @dataclass(frozen=True) class StreamBlockSolveResult: diagnostics: tuple[AlgebraicSolveDiagnostics, ...] scopes: tuple[tuple[str, ...], ...] used_global_fallback: bool @dataclass(frozen=True) class _BlockSolveAttempt: diagnostics: AlgebraicSolveDiagnostics | None optimizer_evaluations: int residual_evaluations: int failure_reason: str | None = None @dataclass(frozen=True) class _MutableAlgebraicStateSnapshot: unknown_values: tuple[tuple[AlgebraicUnknown, float], ...] causal_attributes: tuple[ tuple[object, tuple[tuple[str, object], ...]], ... ] last_diagnostics: AlgebraicSolveDiagnostics | None @classmethod def capture( cls, solver: PressureFlowSolver, ) -> _MutableAlgebraicStateSnapshot: causal_attribute_names = ( "_causal_penetration", "_causal_contact_force", "_causal_port_1_x", "_causal_port_2_x", "_causal_port_1_v", "_causal_port_2_v", ) def component_causal_attributes( component: object, ) -> tuple[tuple[str, object], ...]: names = dict.fromkeys( [ name for name in getattr(component, "__dict__", {}) if name.startswith("_causal_") ] + [ name for name in causal_attribute_names if hasattr(component, name) ] ) return tuple((name, getattr(component, name)) for name in names) return cls( unknown_values=tuple( (unknown, unknown.read()) for unknown in solver.unknowns ), causal_attributes=tuple( ( component, component_causal_attributes(component), ) for component in solver._causal_contact_components ), last_diagnostics=solver.last_diagnostics, ) def restore(self, solver: PressureFlowSolver) -> None: for unknown, value in self.unknown_values: unknown.write(value) for component, attributes in self.causal_attributes: original_names = {name for name, _value in attributes} for name in tuple(getattr(component, "__dict__", {})): if name.startswith("_causal_") and name not in original_names: delattr(component, name) for name, value in attributes: setattr(component, name, value) solver.last_diagnostics = self.last_diagnostics @dataclass(frozen=True) class _ScopedComponentEvaluation: evaluate: Callable[[], tuple[float, ...]] targets: tuple[tuple[int, int, str], ...] @dataclass(frozen=True) class _ScopedConnectionEvaluation: target: int evaluate: Callable[[], float] equation_id: str @dataclass(frozen=True) class _EquationScaleSpec: direct_scale: str | None variable_names: tuple[str, ...] def _scoped_equation_values( equation_count: int, component_evaluations: tuple[_ScopedComponentEvaluation, ...], connection_evaluations: tuple[_ScopedConnectionEvaluation, ...], zero_equation_ids: frozenset[str] = frozenset(), ) -> tuple[float, ...]: values: list[float | None] = [None] * equation_count for evaluation in component_evaluations: if all( equation_id in zero_equation_ids for _target, _source, equation_id in evaluation.targets ): for target, _source, _equation_id in evaluation.targets: values[target] = 0.0 continue component_values = evaluation.evaluate() for target, source, equation_id in evaluation.targets: if equation_id in zero_equation_ids: values[target] = 0.0 continue if source >= len(component_values): raise RuntimeError( "Compiled algebraic equation disappeared at runtime: " f"{equation_id}." ) values[target] = float(component_values[source]) for evaluation in connection_evaluations: values[evaluation.target] = ( 0.0 if evaluation.equation_id in zero_equation_ids else float(evaluation.evaluate()) ) if any(value is None for value in values): raise RuntimeError("Scoped algebraic evaluation returned no value.") return tuple(float(value) for value in values) @dataclass(frozen=True) class _StreamAlgebraicBlock: unknowns: tuple[AlgebraicUnknown, ...] equations: tuple[EquationResidual, ...] component_evaluations: tuple[_ScopedComponentEvaluation, ...] connection_evaluations: tuple[_ScopedConnectionEvaluation, ...] explicit_flow_plan: tuple[ExplicitFlowStage, ...] scope_components: tuple[str, ...] jacobian_entries: tuple[tuple[int, int], ...] equation_scale_specs: tuple[_EquationScaleSpec, ...] def equation_values( self, zero_equation_ids: frozenset[str] = frozenset(), ) -> tuple[float, ...]: return _scoped_equation_values( len(self.equations), self.component_evaluations, self.connection_evaluations, zero_equation_ids, ) @dataclass(frozen=True) class _SelectedEquationEvaluation: equations: tuple[EquationResidual, ...] component_evaluations: tuple[_ScopedComponentEvaluation, ...] connection_evaluations: tuple[_ScopedConnectionEvaluation, ...] block_ranges: tuple[tuple[int, int], ...] equation_scale_specs: tuple[_EquationScaleSpec, ...] scope_components: tuple[str, ...] def equation_values( self, zero_equation_ids: frozenset[str] = frozenset(), ) -> tuple[float, ...]: return _scoped_equation_values( len(self.equations), self.component_evaluations, self.connection_evaluations, zero_equation_ids, ) class StreamPressureBlockSolver: """Re-close only equation blocks whose flow laws consume stream values. The primary pressure-flow solve remains global and performs all mechanical contact and effort causalization. Stream propagation can only invalidate equations owned by components that explicitly declare a stream dependency. For trusted built-in components, the declared ``EquationResidual.variables`` graph identifies the complete square blocks that must be revisited. Any structural ambiguity keeps the legacy full-scope secondary solve. """ def __init__( self, pressure_flow_solver: PressureFlowSolver, sensitive_components: tuple[str, ...], ) -> None: self.pressure_flow_solver = pressure_flow_solver self.sensitive_components = sensitive_components self.fallback_reason: str | None = None self.blocks = self._build_blocks() self._selected_unknowns = tuple( unknown for block in self.blocks for unknown in block.unknowns ) self._selected_unknown_ids = frozenset( unknown.id for unknown in self._selected_unknowns ) self._selected_flow_unknowns = tuple( unknown for block in self.blocks for unknown in block.unknowns if unknown.variable == "m_flow" ) self._selected_explicit_flow_plan = tuple( pressure_flow_solver._compile_explicit_flow_stage( tuple( assignment for assignment in stage.assignments if assignment.unknown.id in self._selected_unknown_ids ) ) for stage in pressure_flow_solver._explicit_flow_plan ) self._selected_equation_evaluation = ( self._compile_selected_equation_evaluation() if self.blocks else None ) ( self._entry_mutated_unknowns, self._unselected_seed_restore_positions, ) = self._compile_mutation_snapshot_plan() ( self._causal_fast_path_eligible, self._causal_fast_path_fallback_reason, self._causal_expected_flow_equation_ids, self._causal_effort_entry_positions, ) = self._build_causal_execution_plan() self._causal_runtime_disabled_reason: str | None = None self._causal_audit_required = True self._causal_solves_since_audit = 0 self._causal_fast_solve_count = 0 self._causal_full_residual_audit_count = 0 self._causal_audit_failure_count = 0 self._causal_legacy_fallback_count = 0 self._causal_last_verified_diagnostics: AlgebraicSolveDiagnostics | None = None @property def available(self) -> bool: return bool(self.blocks) and self.fallback_reason is None @property def causal_fast_path_enabled(self) -> bool: return ( self._causal_fast_path_eligible and self._causal_runtime_disabled_reason is None and self.pressure_flow_solver.causal_fast_path_enabled ) def request_causal_audit(self) -> None: self._causal_audit_required = True def causal_execution_diagnostics(self) -> dict[str, object]: 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") elif not self._causal_fast_path_eligible: disabled_reason = self._causal_fast_path_fallback_reason verified = self._causal_last_verified_diagnostics return { "eligible": self._causal_fast_path_eligible, "enabled": self.causal_fast_path_enabled, "fallbackReason": self._causal_fast_path_fallback_reason, "disabledReason": disabled_reason, "fastSolveCount": self._causal_fast_solve_count, "fullResidualAuditCount": self._causal_full_residual_audit_count, "auditFailureCount": self._causal_audit_failure_count, "legacyFallbackCount": self._causal_legacy_fallback_count, "auditInterval": self.pressure_flow_solver._causal_audit_interval, "solvesSinceAudit": self._causal_solves_since_audit, "lastVerifiedMaxScaledResidual": ( verified.max_scaled_residual if verified is not None else None ), } def _disable_causal_fast_path(self, reason: str) -> None: if self._causal_runtime_disabled_reason is None: self._causal_runtime_disabled_reason = reason self._causal_audit_required = True def _build_causal_execution_plan( self, ) -> tuple[ bool, str | None, frozenset[str], tuple[tuple[AlgebraicUnknown, int], ...], ]: def failed(reason: str): return False, reason, frozenset(), () solver = self.pressure_flow_solver if not self.available: return failed(self.fallback_reason or "streamBlockUnavailable") if not solver.causal_fast_path_eligible: parent_reason = solver.causal_execution_diagnostics()["fallbackReason"] return failed(str(parent_reason or "parentCausalPathIneligible")) if solver._closed_resistance_pressure_plan: return failed("specialClosedResistancePressureSeed") if solver._pnor_pnl0001_series_plan: return failed("specialSeriesPressureSeed") selected = self._selected_equation_evaluation if selected is None: return failed("missingSelectedEquationEvaluation") if any( unknown.variable not in {"p", "m_flow"} for unknown in self._selected_unknowns ): return failed("unsupportedSelectedUnknown") assignments = tuple( assignment for stage in self._selected_explicit_flow_plan for assignment in stage.assignments ) assignment_unknown_ids = tuple( assignment.unknown.id for assignment in assignments ) selected_flow_unknown_ids = frozenset( unknown.id for unknown in self._selected_flow_unknowns ) if ( frozenset(assignment_unknown_ids) != selected_flow_unknown_ids or len(assignment_unknown_ids) != len(selected_flow_unknown_ids) ): return failed("incompleteSelectedFlowCoverage") assignment_equation_ids = tuple( assignment.equation_id for assignment in assignments ) flow_equation_ids = frozenset( equation.id for equation in selected.equations if equation.role == "flow" ) if ( frozenset(assignment_equation_ids) != flow_equation_ids or len(assignment_equation_ids) != len(flow_equation_ids) ): return failed("incompleteSelectedFlowEquationCoverage") if any( equation.role != "effort" for equation in selected.equations if equation.id not in flow_equation_ids ): return failed("selectedResidualIsNotEffortOnly") entry_position_by_id = { unknown.id: position for position, unknown in enumerate(self._entry_mutated_unknowns) } effort_positions: list[tuple[AlgebraicUnknown, int]] = [] for unknown in self._selected_unknowns: if unknown.variable != "p": continue position = entry_position_by_id.get(unknown.id) if position is None: return failed("missingEffortMutationSnapshot") effort_positions.append((unknown, position)) return ( True, None, flow_equation_ids, tuple(effort_positions), ) def _causal_audit_is_due(self) -> bool: return ( self._causal_audit_required or self._causal_last_verified_diagnostics is None or self._causal_solves_since_audit >= self.pressure_flow_solver._causal_audit_interval ) def _record_causal_audit( self, diagnostics: AlgebraicSolveDiagnostics, ) -> None: self._causal_full_residual_audit_count += 1 self._causal_solves_since_audit = 0 self._causal_audit_required = False self._causal_last_verified_diagnostics = diagnostics def _causal_fast_diagnostics( self, scale_context: Mapping[str, float], ) -> AlgebraicSolveDiagnostics: verified = self._causal_last_verified_diagnostics if verified is None: raise RuntimeError("Stream causal execution has no residual audit.") return replace( verified, message=( "Compiled causal stream-pressure block completed; residuals " "reuse the latest full audit." ), evaluations=0, pressure_scale=float(scale_context["p"]), flow_scale=float(scale_context["m_flow"]), 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, ) @staticmethod def _equation_scale_spec(equation: EquationResidual) -> _EquationScaleSpec: variable_names = tuple( variable.rsplit(".", 1)[-1] for variable in equation.variables ) if equation.role == "flow": return _EquationScaleSpec("m_flow", variable_names) if equation.role == "effort" or "p" in variable_names: return _EquationScaleSpec("p", variable_names) return _EquationScaleSpec(None, variable_names) def _build_blocks(self) -> tuple[_StreamAlgebraicBlock, ...]: solver = self.pressure_flow_solver if not solver.jacobian_sparsity_is_trusted: self.fallback_reason = ( solver.jacobian_sparsity_fallback_reason or "untrustedAlgebraicStructure" ) return () unknowns = solver.unknowns equations = solver.equation_templates unknown_index_by_id = { unknown.id: index for index, unknown in enumerate(unknowns) } equation_unknowns: list[tuple[int, ...]] = [] equations_by_unknown: list[list[int]] = [[] for _unknown in unknowns] for equation_index, equation in enumerate(equations): dependencies = tuple( dict.fromkeys( unknown_index_by_id[variable] for variable in equation.variables if variable in unknown_index_by_id ) ) if not dependencies: self.fallback_reason = "equationWithoutDeclaredUnknown" return () equation_unknowns.append(dependencies) for unknown_index in dependencies: equations_by_unknown[unknown_index].append(equation_index) if any(not attached for attached in equations_by_unknown): self.fallback_reason = "unknownWithoutDeclaredEquation" return () raw_blocks: list[tuple[tuple[int, ...], tuple[int, ...]]] = [] visited_unknowns: set[int] = set() visited_equations: set[int] = set() for root_unknown in range(len(unknowns)): if root_unknown in visited_unknowns: continue block_unknowns: set[int] = set() block_equations: set[int] = set() pending_unknowns = [root_unknown] while pending_unknowns: unknown_index = pending_unknowns.pop() if unknown_index in visited_unknowns: continue visited_unknowns.add(unknown_index) block_unknowns.add(unknown_index) for equation_index in equations_by_unknown[unknown_index]: if equation_index not in visited_equations: visited_equations.add(equation_index) block_equations.add(equation_index) for dependency in equation_unknowns[equation_index]: if dependency not in visited_unknowns: pending_unknowns.append(dependency) ordered_unknowns = tuple(sorted(block_unknowns)) ordered_equations = tuple(sorted(block_equations)) if len(ordered_unknowns) != len(ordered_equations): self.fallback_reason = "nonSquareEquationBlock" return () raw_blocks.append((ordered_unknowns, ordered_equations)) if len(visited_equations) != len(equations): self.fallback_reason = "unreachableEquationBlock" return () sensitive = frozenset(self.sensitive_components) if sensitive - set(solver.network.components): self.fallback_reason = "unknownSensitiveComponent" return () component_equation_owners = { equation.owner_id for equation in equations if equation.owner == "component" } if sensitive - component_equation_owners: self.fallback_reason = "sensitiveComponentHasNoEquationBlock" return () selected = [ block for block in raw_blocks if any( equations[equation_index].owner == "component" and equations[equation_index].owner_id in sensitive for equation_index in block[1] ) ] if not selected: self.fallback_reason = "sensitiveComponentHasNoEquationBlock" return () for unknown_indices, equation_indices in selected: if { unknowns[unknown_index].variable for unknown_index in unknown_indices } - {"p", "m_flow"}: self.fallback_reason = "streamBlockContainsNonPneumaticUnknown" return () component_names = { unknowns[unknown_index].component for unknown_index in unknown_indices } | { equations[equation_index].owner_id for equation_index in equation_indices if equations[equation_index].owner == "component" } if component_names - set(solver.network.components): self.fallback_reason = "unknownEquationOwner" return () if any( not type(solver.network.components[name]).__module__.startswith( "app.simulation.components." ) for name in component_names ): self.fallback_reason = "untrustedCustomComponent" return () equation_locations: list[tuple[str, object, int]] = [] for component_plan in solver._component_equation_plan: equation_locations.extend( ("component", component_plan, local_index) for local_index, _template in enumerate(component_plan.templates) ) equation_locations.extend( ("connection", connection_plan, 0) for connection_plan in solver._connection_equation_plan ) if len(equation_locations) != len(equations): self.fallback_reason = "equationEvaluationPlanMismatch" return () compiled: list[_StreamAlgebraicBlock] = [] for unknown_indices, equation_indices in selected: block_unknowns = tuple(unknowns[index] for index in unknown_indices) block_equations = tuple(equations[index] for index in equation_indices) block_unknown_ids = frozenset(unknown.id for unknown in block_unknowns) component_targets: dict[int, list[tuple[int, int, str]]] = {} component_plans: dict[int, object] = {} connection_evaluations: list[_ScopedConnectionEvaluation] = [] for target, global_equation_index in enumerate(equation_indices): kind, evaluation_plan, source = equation_locations[ global_equation_index ] if kind == "component": key = id(evaluation_plan) component_plans[key] = evaluation_plan component_targets.setdefault(key, []).append( ( target, source, equations[global_equation_index].id, ) ) else: connection_evaluations.append( _ScopedConnectionEvaluation( target=target, evaluate=evaluation_plan.evaluate, equation_id=equations[global_equation_index].id, ) ) component_evaluations = tuple( _ScopedComponentEvaluation( evaluate=component_plans[key].evaluate, targets=tuple(targets), ) for key, targets in component_targets.items() ) explicit_flow_plan = tuple( solver._compile_explicit_flow_stage( tuple( assignment for assignment in stage.assignments if assignment.unknown.id in block_unknown_ids ) ) for stage in solver._explicit_flow_plan ) local_unknown_index = { unknown.id: index for index, unknown in enumerate(block_unknowns) } jacobian_entries = tuple( (row, local_unknown_index[variable]) for row, equation in enumerate(block_equations) for variable in dict.fromkeys(equation.variables) if variable in local_unknown_index ) scope_components = tuple( dict.fromkeys( [unknown.component for unknown in block_unknowns] + [ equation.owner_id for equation in block_equations if equation.owner == "component" ] ) ) compiled.append( _StreamAlgebraicBlock( unknowns=block_unknowns, equations=block_equations, component_evaluations=component_evaluations, connection_evaluations=tuple(connection_evaluations), explicit_flow_plan=explicit_flow_plan, scope_components=scope_components, jacobian_entries=jacobian_entries, equation_scale_specs=tuple( self._equation_scale_spec(equation) for equation in block_equations ), ) ) return tuple(compiled) def _compile_selected_equation_evaluation( self, ) -> _SelectedEquationEvaluation: equations: list[EquationResidual] = [] equation_scale_specs: list[_EquationScaleSpec] = [] component_targets: dict[int, list[tuple[int, int, str]]] = {} component_evaluators: dict[int, Callable[[], tuple[float, ...]]] = {} connection_evaluations: list[_ScopedConnectionEvaluation] = [] block_ranges: list[tuple[int, int]] = [] scope_components: list[str] = [] offset = 0 for block in self.blocks: start = offset equations.extend(block.equations) equation_scale_specs.extend(block.equation_scale_specs) scope_components.extend(block.scope_components) for evaluation in block.component_evaluations: key = id(evaluation.evaluate) component_evaluators[key] = evaluation.evaluate component_targets.setdefault(key, []).extend( (offset + target, source, equation_id) for target, source, equation_id in evaluation.targets ) connection_evaluations.extend( _ScopedConnectionEvaluation( target=offset + evaluation.target, evaluate=evaluation.evaluate, equation_id=evaluation.equation_id, ) for evaluation in block.connection_evaluations ) offset += len(block.equations) block_ranges.append((start, offset)) return _SelectedEquationEvaluation( equations=tuple(equations), component_evaluations=tuple( _ScopedComponentEvaluation( evaluate=component_evaluators[key], targets=tuple(targets), ) for key, targets in component_targets.items() ), connection_evaluations=tuple(connection_evaluations), block_ranges=tuple(block_ranges), equation_scale_specs=tuple(equation_scale_specs), scope_components=tuple(dict.fromkeys(scope_components)), ) def _compile_mutation_snapshot_plan( self, ) -> tuple[ tuple[AlgebraicUnknown, ...], tuple[tuple[AlgebraicUnknown, int], ...], ]: if not self.blocks: return (), () solver = self.pressure_flow_solver unknowns_by_id = {unknown.id: unknown for unknown in solver.unknowns} special_seed_target_ids: list[str] = [] for binding in solver._closed_resistance_pressure_plan: special_seed_target_ids.extend( ( f"{binding.component.name}.{binding.port_name}.p", f"{binding.neighbor.name}.{binding.neighbor_port}.p", ) ) for binding in solver._pnor_pnl0001_series_plan: special_seed_target_ids.extend( ( f"{binding.orifice.name}.{binding.orifice_port}.p", f"{binding.pipe.name}.{binding.pipe_port}.p", ) ) entry_mutated: list[AlgebraicUnknown] = [] seen_ids: set[str] = set() def append_unknown(unknown: AlgebraicUnknown | None) -> None: if unknown is None or unknown.id in seen_ids: return seen_ids.add(unknown.id) entry_mutated.append(unknown) for block in self.blocks: for unknown in block.unknowns: append_unknown(unknown) for target_id in special_seed_target_ids: append_unknown(unknowns_by_id.get(target_id)) positions = { unknown.id: index for index, unknown in enumerate(entry_mutated) } unselected_seed_restore_positions = tuple( (unknowns_by_id[target_id], positions[target_id]) for target_id in dict.fromkeys(special_seed_target_ids) if target_id in unknowns_by_id and target_id not in self._selected_unknown_ids ) return tuple(entry_mutated), unselected_seed_restore_positions def _seed_selected_blocks( self, entry_values: tuple[float, ...], ) -> frozenset[str]: solver = self.pressure_flow_solver seeded_equation_ids: set[str] = set() try: # Reuse the global special seed plans because they encode catalog # behavior such as PNOR/PNL0001 series pressure initialization. # Shared PortState objects make those helpers capable of touching # other equation blocks, so their exact unselected pressure targets # are restored before block residuals are evaluated. Secondary # closure deliberately mirrors ``solver.solve(effort_variables=())``: # the primary global solve has already propagated equal pressures. solver._seed_closed_resistance_pressures() solver._seed_pnor_pnl0001_series_pressures() for unknown in self._selected_flow_unknowns: unknown.write(0.0) for stage in self._selected_explicit_flow_plan: values = solver._evaluate_explicit_flow_stage(stage) for assignment, target_value in zip(stage.assignments, values): if isfinite(target_value): assignment.unknown.write(target_value) seeded_equation_ids.add(assignment.equation_id) finally: for unknown, position in self._unselected_seed_restore_positions: unknown.write(entry_values[position]) return frozenset(seeded_equation_ids) @staticmethod def _equation_scales_from_specs( specs: tuple[_EquationScaleSpec, ...], scales: Mapping[str, float], ) -> tuple[float, ...]: result: list[float] = [] for spec in specs: if spec.direct_scale is not None: result.append(float(scales[spec.direct_scale])) else: result.append( max( [ float(scales.get(name, 1.0)) for name in spec.variable_names ] + [1.0] ) ) return tuple(result) @classmethod def _equation_scales( cls, block: _StreamAlgebraicBlock, scales: Mapping[str, float], ) -> tuple[float, ...]: return cls._equation_scales_from_specs( block.equation_scale_specs, scales, ) def _seeded_diagnostics( self, *, unknowns: tuple[AlgebraicUnknown, ...], equation_values: tuple[float, ...], equation_scales: tuple[float, ...], scale_context: Mapping[str, float], message: str, ) -> AlgebraicSolveDiagnostics | None: scaled = tuple( abs(value / scale) for value, scale in zip(equation_values, equation_scales) ) feasible = all( isfinite(unknown.read()) and ( unknown.variable != "p" or unknown.read() > PRESSURE_LOWER_BOUND_PA ) for unknown in unknowns ) max_scaled = max(scaled, default=0.0) if ( not feasible or not all(isfinite(value) for value in scaled) or max_scaled > self.pressure_flow_solver.residual_tolerance ): return None return AlgebraicSolveDiagnostics( success=True, message=message, evaluations=0, pressure_scale=float(scale_context["p"]), flow_scale=float(scale_context["m_flow"]), max_scaled_residual=max_scaled, max_raw_residual=max( (abs(value) for value in equation_values), default=0.0, ), ) def _solve_block( self, block: _StreamAlgebraicBlock, scale_context: Mapping[str, float], *, seeded_values: tuple[float, ...] | None = None, equation_scales: tuple[float, ...] | None = None, ) -> _BlockSolveAttempt: import numpy as np from scipy.optimize import least_squares from scipy.sparse import csr_matrix if equation_scales is None: equation_scales = self._equation_scales(block, scale_context) if seeded_values is None: seeded_values = block.equation_values() solver = self.pressure_flow_solver seeded_diagnostics = self._seeded_diagnostics( unknowns=block.unknowns, equation_values=seeded_values, equation_scales=equation_scales, scale_context=scale_context, message=( "Seeded stream-sensitive algebraic block satisfies the " "residual tolerance." ), ) if seeded_diagnostics is not None: return _BlockSolveAttempt( diagnostics=seeded_diagnostics, optimizer_evaluations=0, residual_evaluations=0, ) unknown_scales = tuple( float(scale_context[unknown.variable]) for unknown in block.unknowns ) fallback_pressure = float(scale_context["fallback_pressure"]) x0 = np.asarray( [ ( unknown.read() if unknown.variable != "p" or unknown.read() > 0.0 else fallback_pressure ) / scale for unknown, scale in zip(block.unknowns, unknown_scales) ], dtype=float, ) lower = np.asarray( [ ( PRESSURE_LOWER_BOUND_PA / float(scale_context["p"]) if unknown.variable == "p" else -np.inf ) for unknown in block.unknowns ] ) upper = np.full(len(block.unknowns), np.inf) def assign(values) -> None: for unknown, value, scale in zip( block.unknowns, values, unknown_scales, ): unknown.write(float(value) * scale) residual_evaluations = 0 def scaled_residuals(values): nonlocal residual_evaluations residual_evaluations += 1 assign(values) return np.asarray( [ value / scale for value, scale in zip( block.equation_values(), equation_scales, ) ], dtype=float, ) rows = [row for row, _column in block.jacobian_entries] columns = [column for _row, column in block.jacobian_entries] jacobian_sparsity = csr_matrix( ( np.ones(len(rows), dtype=bool), (rows, columns), ), shape=(len(block.equations), len(block.unknowns)), ) if ( any(jacobian_sparsity.getnnz(axis=1) == 0) or any(jacobian_sparsity.getnnz(axis=0) == 0) ): return _BlockSolveAttempt( diagnostics=None, optimizer_evaluations=0, residual_evaluations=0, failure_reason="invalidBlockJacobianSparsity", ) result = None try: result = least_squares( scaled_residuals, x0, bounds=(lower, upper), jac_sparsity=jacobian_sparsity, x_scale="jac", ftol=1.0e-10, xtol=1.0e-10, gtol=1.0e-10, max_nfev=solver.max_evaluations, ) assign(result.x) equation_values = block.equation_values() except MemoryError: assign(x0) raise except Exception as exc: assign(x0) return _BlockSolveAttempt( diagnostics=None, optimizer_evaluations=( int(result.nfev) if result is not None else 0 ), residual_evaluations=residual_evaluations, failure_reason=f"blockSolveFailed:{type(exc).__name__}", ) except BaseException: assign(x0) raise scaled = tuple( abs(value / scale) for value, scale in zip(equation_values, equation_scales) ) max_scaled = max(scaled, default=0.0) success = ( all(isfinite(value) for value in scaled) and max_scaled <= solver.residual_tolerance and (bool(result.success) or int(result.status) == 0) ) if not success: assign(x0) return _BlockSolveAttempt( diagnostics=None, optimizer_evaluations=int(result.nfev), residual_evaluations=residual_evaluations, failure_reason="blockResidualNotConverged", ) diagnostics = AlgebraicSolveDiagnostics( success=True, message=str(result.message), evaluations=int(result.nfev), pressure_scale=float(scale_context["p"]), flow_scale=float(scale_context["m_flow"]), max_scaled_residual=max_scaled, max_raw_residual=max( (abs(value) for value in equation_values), default=0.0, ), residual_evaluations=residual_evaluations, jacobian_mode="sparse", dense_fallback_used=False, nonlinear_block_count=1, nonlinear_block_unknown_count=len(block.unknowns), ) return _BlockSolveAttempt( diagnostics=diagnostics, optimizer_evaluations=int(result.nfev), residual_evaluations=residual_evaluations, ) @staticmethod def _aggregate_local_diagnostics( diagnostics: tuple[AlgebraicSolveDiagnostics, ...], *, nonlinear_block_unknown_count: int, ) -> AlgebraicSolveDiagnostics: if len(diagnostics) == 1: return diagnostics[0] nonlinear = tuple( item for item in diagnostics if item.jacobian_mode != "seeded" ) representative = diagnostics[-1] return replace( representative, message=( f"{len(diagnostics)} stream-sensitive algebraic blocks " "satisfied the residual tolerance." ), evaluations=sum(item.evaluations for item in diagnostics), max_scaled_residual=max( item.max_scaled_residual for item in diagnostics ), max_raw_residual=max(item.max_raw_residual for item in diagnostics), residual_evaluations=sum( item.residual_evaluations for item in diagnostics ), jacobian_mode="blockSparse" if nonlinear else "seeded", dense_fallback_used=any( item.dense_fallback_used for item in diagnostics ), nonlinear_block_count=len(nonlinear), nonlinear_block_unknown_count=( nonlinear_block_unknown_count if nonlinear else 0 ), block_fallback_used=False, block_fallback_reason=None, ) def solve( self, *, 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 ) entry_values = tuple( unknown.read() for unknown in self._entry_mutated_unknowns ) def restore_entry_mutations() -> None: for unknown, value in zip( self._entry_mutated_unknowns, entry_values, ): unknown.write(value) def global_fallback( *, reason: str | None, local_optimizer_evaluations: int = 0, local_residual_evaluations: int = 0, local_attempted: bool, ) -> StreamBlockSolveResult: restore_entry_mutations() fallback_snapshot = _MutableAlgebraicStateSnapshot.capture(solver) try: fallback = solver.solve( effort_variables=(), scale_context=context, ) except BaseException: fallback_snapshot.restore(solver) raise fallback_reason = reason or fallback.block_fallback_reason if ( reason and fallback.block_fallback_reason and reason != fallback.block_fallback_reason ): fallback_reason = ( f"{reason};global:{fallback.block_fallback_reason}" ) aggregate = replace( fallback, evaluations=( local_optimizer_evaluations + fallback.evaluations ), residual_evaluations=( local_residual_evaluations + fallback.residual_evaluations ), block_fallback_used=( local_attempted or fallback.block_fallback_used ), block_fallback_reason=fallback_reason, ) return StreamBlockSolveResult( diagnostics=(aggregate,), scopes=(tuple(solver.network.components),), used_global_fallback=True, ) if not self.available: return global_fallback( reason=self.fallback_reason, local_attempted=False, ) try: seeded_equation_ids = ( self._seed_selected_blocks(entry_values) or frozenset() ) except MemoryError: restore_entry_mutations() raise except Exception as exc: if causal_candidate: self._causal_legacy_fallback_count += 1 self._disable_causal_fast_path( f"causalSecondarySeedFailed:{type(exc).__name__}" ) return global_fallback( reason=f"secondarySeedFailed:{type(exc).__name__}", local_attempted=True, ) except BaseException: restore_entry_mutations() raise if causal_candidate: causal_runtime_gate_passed = ( seeded_equation_ids == self._causal_expected_flow_equation_ids and all( isfinite(unknown.read()) for unknown in self._selected_flow_unknowns ) and all( unknown.read() == entry_values[position] for unknown, position in self._causal_effort_entry_positions ) ) if not causal_runtime_gate_passed: self._causal_legacy_fallback_count += 1 self._disable_causal_fast_path("causalSecondaryRuntimeGateFailed") causal_candidate = False causal_audit_due = False elif not causal_audit_due: diagnostics = self._causal_fast_diagnostics(context) self._causal_fast_solve_count += 1 self._causal_solves_since_audit += 1 return StreamBlockSolveResult( diagnostics=(diagnostics,), scopes=( self._selected_equation_evaluation.scope_components, ), used_global_fallback=False, ) selected_evaluation = self._selected_equation_evaluation assert selected_evaluation is not None try: selected_values = selected_evaluation.equation_values( seeded_equation_ids ) selected_scales = self._equation_scales_from_specs( selected_evaluation.equation_scale_specs, context, ) union_diagnostics = self._seeded_diagnostics( unknowns=self._selected_unknowns, equation_values=selected_values, equation_scales=selected_scales, scale_context=context, message=( "Seeded stream-sensitive algebraic equation blocks satisfy " "the residual tolerance." ), ) except MemoryError: restore_entry_mutations() raise except Exception as exc: if causal_candidate: self._causal_legacy_fallback_count += 1 self._disable_causal_fast_path( "causalSecondaryResidualEvaluationFailed:" f"{type(exc).__name__}" ) return global_fallback( reason=( "secondaryResidualEvaluationFailed:" f"{type(exc).__name__}" ), local_attempted=True, ) except BaseException: restore_entry_mutations() raise if union_diagnostics is not None: if causal_candidate and causal_audit_due: self._record_causal_audit(union_diagnostics) return StreamBlockSolveResult( diagnostics=(union_diagnostics,), scopes=(selected_evaluation.scope_components,), used_global_fallback=False, ) if causal_candidate and causal_audit_due: self._causal_audit_failure_count += 1 self._causal_legacy_fallback_count += 1 self._disable_causal_fast_path("causalSecondaryResidualAuditFailed") causal_candidate = False diagnostics: list[AlgebraicSolveDiagnostics] = [] local_optimizer_evaluations = 0 local_residual_evaluations = 0 nonlinear_block_unknown_count = 0 for block, (start, stop) in zip( self.blocks, selected_evaluation.block_ranges, ): try: attempt = self._solve_block( block, context, seeded_values=selected_values[start:stop], equation_scales=selected_scales[start:stop], ) except MemoryError: restore_entry_mutations() raise except Exception as exc: return global_fallback( reason=f"secondaryBlockSolveFailed:{type(exc).__name__}", local_optimizer_evaluations=local_optimizer_evaluations, local_residual_evaluations=local_residual_evaluations, local_attempted=True, ) except BaseException: restore_entry_mutations() raise local_optimizer_evaluations += attempt.optimizer_evaluations local_residual_evaluations += attempt.residual_evaluations if attempt.diagnostics is None: return global_fallback( reason=attempt.failure_reason or "secondaryBlockSolveFailed", local_optimizer_evaluations=local_optimizer_evaluations, local_residual_evaluations=local_residual_evaluations, local_attempted=True, ) diagnostics.append(attempt.diagnostics) if attempt.diagnostics.jacobian_mode != "seeded": nonlinear_block_unknown_count += len(block.unknowns) aggregate = self._aggregate_local_diagnostics( tuple(diagnostics), nonlinear_block_unknown_count=nonlinear_block_unknown_count, ) return StreamBlockSolveResult( diagnostics=(aggregate,), scopes=(selected_evaluation.scope_components,), used_global_fallback=False, )