优化原生结果编码传输与浏览器缓存,记录八路性能基线
原生结果series通过字节索引直传,C端使用Ryu精确回读编码和64 KiB批量写出;网页采用Float64缓存和CSV工作线程,减少结果处理与保存等待。 补充八路AME曲线核查、全流程分阶段计时、独立编码基准和复现工具,固定后续优化采用修正八路及rtol=1e-8。C写出1.1808→0.1638 s,点击到可查看8.0100→6.9756 s。 验证:最终10项编码专项、29项相关后端回归通过;8份原生结果逐位一致,16次网页结果/CSV/刷新恢复通过。前端构建及缓存/CSV专项在本轮结果处理工作中通过。环境、原始大结果与临时构建不纳入Git。
This commit is contained in:
1 parent
808c484f5b
commit
3bc4be3c06
55 files changed
+20522
-141
No files matched your search
+22
-7
@@ -23,6 +23,7 @@ from fastapi.responses import FileResponse, HTMLResponse, StreamingResponse
|
||||
from pydantic import BaseModel, ConfigDict, Field, ValidationError
|
||||
|
||||
from app.simulation.performance import performance_span, profile_phase, profile_run
|
||||
from app.simulation.native_codegen.transport import NativeSeriesJson, serialize_result_parts
|
||||
from app.simulation.config import SolverActivityTracker
|
||||
from app.system_xml import (
|
||||
SystemXmlDocument,
|
||||
@@ -652,7 +653,7 @@ async def simulate_system_xml_stream(request: Request) -> StreamingResponse:
|
||||
simulation_id = request.headers.get("x-simulation-id") or uuid4().hex
|
||||
task = _register_simulation_task(simulation_id)
|
||||
return StreamingResponse(
|
||||
simulation_event_stream(await request.body(), task=task),
|
||||
simulation_event_stream(await request.body(), task=task, raw_series=True),
|
||||
media_type="application/x-ndjson",
|
||||
headers={
|
||||
"Cache-Control": "no-cache, no-transform",
|
||||
@@ -679,13 +680,17 @@ def cancel_system_xml_simulation(
|
||||
}
|
||||
|
||||
|
||||
@app.get("/api/system-xml/simulations/{simulation_id}")
|
||||
def get_system_xml_simulation(simulation_id: str) -> dict[str, object]:
|
||||
@app.get("/api/system-xml/simulations/{simulation_id}", response_model=None)
|
||||
def get_system_xml_simulation(simulation_id: str) -> dict[str, object] | Response:
|
||||
with SIMULATION_TASKS_LOCK:
|
||||
task = SIMULATION_TASKS.get(simulation_id)
|
||||
if task is None:
|
||||
raise HTTPException(status_code=404, detail="Simulation task was not found.")
|
||||
return _simulation_task_snapshot(task)
|
||||
snapshot = _simulation_task_snapshot(task)
|
||||
result = snapshot.get("result")
|
||||
if isinstance(result, dict) and isinstance(result.get("series"), NativeSeriesJson):
|
||||
return StreamingResponse(iter(serialize_result_parts(snapshot)), media_type="application/json")
|
||||
return snapshot
|
||||
|
||||
|
||||
def run_system_xml_simulation(
|
||||
@@ -693,6 +698,7 @@ def run_system_xml_simulation(
|
||||
progress_callback: SimulationProgressEmitter | None = None,
|
||||
cancel_check: Callable[[], bool] | None = None,
|
||||
activity_tracker: SolverActivityTracker | None = None,
|
||||
*, raw_series: bool = False,
|
||||
) -> dict[str, object]:
|
||||
with profile_run() as trace:
|
||||
result = _run_system_xml_simulation_profiled(
|
||||
@@ -700,6 +706,7 @@ def run_system_xml_simulation(
|
||||
progress_callback,
|
||||
cancel_check,
|
||||
activity_tracker,
|
||||
raw_series=raw_series,
|
||||
)
|
||||
|
||||
performance = trace.snapshot()
|
||||
@@ -715,6 +722,7 @@ def _run_system_xml_simulation_profiled(
|
||||
progress_callback: SimulationProgressEmitter | None = None,
|
||||
cancel_check: Callable[[], bool] | None = None,
|
||||
activity_tracker: SolverActivityTracker | None = None,
|
||||
*, raw_series: bool = False,
|
||||
) -> dict[str, object]:
|
||||
from app.simulation.backends import simulate_network
|
||||
from app.simulation.results import SimulationPreparationError
|
||||
@@ -762,6 +770,7 @@ def _run_system_xml_simulation_profiled(
|
||||
progress_callback=report_system_progress,
|
||||
cancel_check=cancel_check,
|
||||
activity_tracker=activity_tracker,
|
||||
raw_series=raw_series,
|
||||
)
|
||||
except SimulationPreparationError as exc:
|
||||
raise HTTPException(
|
||||
@@ -799,7 +808,7 @@ def _run_system_xml_simulation_profiled(
|
||||
"validation": report.as_dict(),
|
||||
"simulation": document.as_model_data()["simulation"],
|
||||
"model": network.as_interface_dict(),
|
||||
**result.as_dict(),
|
||||
**result.as_dict(raw_series=raw_series),
|
||||
}
|
||||
|
||||
|
||||
@@ -807,7 +816,8 @@ def simulation_event_stream(
|
||||
xml_bytes: bytes,
|
||||
*,
|
||||
task: SimulationTaskRecord | None = None,
|
||||
) -> Iterator[str]:
|
||||
raw_series: bool = False,
|
||||
) -> Iterator[str | bytes]:
|
||||
events: queue.Queue[dict[str, object] | object] = queue.Queue()
|
||||
finished = object()
|
||||
latest_progress = 0
|
||||
@@ -857,6 +867,7 @@ def simulation_event_stream(
|
||||
emit_progress,
|
||||
task.cancel_event.is_set if task is not None else None,
|
||||
activity_tracker,
|
||||
**({"raw_series": True} if raw_series else {}),
|
||||
)
|
||||
if task is not None:
|
||||
result = _mark_simulation_task_result(task, result)
|
||||
@@ -958,7 +969,11 @@ def simulation_event_stream(
|
||||
continue
|
||||
if event is finished:
|
||||
break
|
||||
yield json.dumps(event, ensure_ascii=False, separators=(",", ":")) + "\n"
|
||||
if raw_series and isinstance(event, dict) and event.get("event") == "result":
|
||||
yield from serialize_result_parts(event)
|
||||
yield b"\n"
|
||||
else:
|
||||
yield json.dumps(event, ensure_ascii=False, separators=(",", ":")) + "\n"
|
||||
finally:
|
||||
if task is not None:
|
||||
_request_simulation_task_cancel(task, "stalled")
|
||||
|
||||
@@ -30,9 +30,9 @@ def simulation_config(simulation) -> SolveIVPConfig:
|
||||
|
||||
|
||||
def simulate_network(network, simulation, *, progress_callback=None,
|
||||
cancel_check=None, activity_tracker=None, backend=None):
|
||||
cancel_check=None, activity_tracker=None, backend=None, raw_series=False):
|
||||
numeric_engine_name(backend)
|
||||
config = simulation_config(simulation)
|
||||
from app.simulation.native_codegen.runner import simulate_native
|
||||
return simulate_native(network, config, sample_step=simulation.sample_step, progress_callback=progress_callback,
|
||||
cancel_check=cancel_check, activity_tracker=activity_tracker)
|
||||
cancel_check=cancel_check, activity_tracker=activity_tracker, raw_series=raw_series)
|
||||
@@ -52,7 +52,7 @@ def toolchain() -> tuple[str, Path, str]:
|
||||
def build_native(program: NativeProgram, *, cache_dir: Path | None = None) -> NativeBuild:
|
||||
start = time.perf_counter()
|
||||
compiler, sundials, compiler_version = toolchain()
|
||||
runtime = sorted(NATIVE.rglob("*.c")) + sorted((NATIVE / "include").glob("*.h"))
|
||||
runtime = sorted(NATIVE.rglob("*.c")) + sorted((NATIVE / "include").rglob("*.h"))
|
||||
flags = ["-std=c11", "-O3", "-Wall", "-Wextra", "-Werror", "-ffp-contract=off", "-fno-fast-math"]
|
||||
executable_name = "model.exe" if os.name == "nt" else "model"
|
||||
if os.name == "nt":
|
||||
|
||||
@@ -12,13 +12,14 @@ import time
|
||||
|
||||
from app.simulation.config import SolveIVPConfig
|
||||
from app.simulation.results import GenericSimulationResult
|
||||
from .transport import NativeSeriesJson, read_indexed_result
|
||||
from .build import NativeBuild, build_native
|
||||
from .compiler import NativeCapabilityError, compile_native_program
|
||||
|
||||
|
||||
def execute_native(build: NativeBuild, config: SolveIVPConfig, sample_step: float, *,
|
||||
run_dir: Path, record_samples=True, cancel_check=None,
|
||||
progress_callback=None, activity_tracker=None, timeout=300.0) -> dict:
|
||||
progress_callback=None, activity_tracker=None, timeout=300.0, raw_series=False) -> dict:
|
||||
if config.method not in ("RK45", "BDF"):
|
||||
raise NativeCapabilityError(f"Native v1 does not support method {config.method}.")
|
||||
if not isinstance(config.atol, (int, float)) or config.atol != 1e-8 or config.first_step is not None:
|
||||
@@ -26,13 +27,16 @@ def execute_native(build: NativeBuild, config: SolveIVPConfig, sample_step: floa
|
||||
run_dir.mkdir(parents=True, exist_ok=True)
|
||||
output = run_dir / "result.json"
|
||||
cancel_path = run_dir / "cancel.request"
|
||||
if output.exists() or cancel_path.exists():
|
||||
index_path = run_dir / "result-index.json"
|
||||
if output.exists() or cancel_path.exists() or index_path.exists():
|
||||
raise ValueError("Native execution requires a fresh run directory.")
|
||||
command = [str(build.executable), "--method", config.method,
|
||||
"--start", str(config.t_start), "--stop", str(config.t_stop),
|
||||
"--sample-step", str(sample_step), "--max-step", str(config.max_step),
|
||||
"--rtol", str(config.rtol), "--timeout", str(timeout),
|
||||
"--cancel-file", str(cancel_path.resolve()), "--output", str(output.resolve())]
|
||||
if raw_series:
|
||||
command.extend(["--result-index", str(index_path.resolve())])
|
||||
if not record_samples:
|
||||
command.append("--solve-only")
|
||||
creationflags = subprocess.CREATE_NO_WINDOW if os.name == "nt" else 0
|
||||
@@ -86,9 +90,13 @@ def execute_native(build: NativeBuild, config: SolveIVPConfig, sample_step: floa
|
||||
process.stderr.close()
|
||||
if not output.is_file():
|
||||
raise RuntimeError(f"Native worker exited with code {process.returncode} without results; see {run_dir / 'worker.log'}.")
|
||||
payload = json.loads(output.read_text(encoding="utf-8"))
|
||||
if process.returncode not in (0, 2):
|
||||
raise RuntimeError(f"Native worker failed with exit code {process.returncode}.")
|
||||
try:
|
||||
payload = (read_indexed_result(output, index_path) if raw_series
|
||||
else json.loads(output.read_text(encoding="utf-8")))
|
||||
except OSError as exc:
|
||||
raise RuntimeError(f"Cannot read native worker result artifacts: {exc}") from exc
|
||||
payload["processWallSeconds"] = time.perf_counter()-started
|
||||
payload["buildKey"] = build.manifest["buildKey"]
|
||||
payload["cacheHit"] = build.cache_hit
|
||||
@@ -99,7 +107,7 @@ def execute_native(build: NativeBuild, config: SolveIVPConfig, sample_step: floa
|
||||
|
||||
|
||||
def simulate_native(network, config, *, sample_step, progress_callback=None,
|
||||
cancel_check=None, activity_tracker=None):
|
||||
cancel_check=None, activity_tracker=None, raw_series=False):
|
||||
if config.method not in ("RK45", "BDF"):
|
||||
raise NativeCapabilityError(f"Native v1 does not support method {config.method}.")
|
||||
if progress_callback:
|
||||
@@ -109,7 +117,7 @@ def simulate_native(network, config, *, sample_step, progress_callback=None,
|
||||
with tempfile.TemporaryDirectory(prefix="native-simulation-") as directory:
|
||||
data = execute_native(build, config, sample_step, run_dir=Path(directory),
|
||||
cancel_check=cancel_check, progress_callback=progress_callback,
|
||||
activity_tracker=activity_tracker)
|
||||
activity_tracker=activity_tracker, raw_series=raw_series)
|
||||
totals = {
|
||||
"nfev": data["nfev"], "njev": data["njev"], "nlu": data["nlu"],
|
||||
"acceptedStepCount": data["acceptedSteps"], "rejectedStepCount": data["rejectedSteps"],
|
||||
@@ -124,5 +132,6 @@ def simulate_native(network, config, *, sample_step, progress_callback=None,
|
||||
variables=program.variables, series=data["series"], final=data["final"],
|
||||
diagnostics={"backend": "native-c", "native": {k: v for k, v in data.items()
|
||||
if k not in ("series", "final", "finalState")}, "integration": {"method": config.method, "rtol": config.rtol, "totals": totals},
|
||||
"stateCount": len(program.state_keys), "sampleCount": len(data["series"]["time"])},
|
||||
"stateCount": len(program.state_keys), "sampleCount": (data["series"].sample_count if isinstance(data["series"], NativeSeriesJson)
|
||||
else len(data["series"].get("time", [])))},
|
||||
)
|
||||
@@ -0,0 +1,63 @@
|
||||
"""Carry trusted C-generated series JSON without a Python float-array round trip.
|
||||
|
||||
Only native output plus its byte index may construct this transport object. Public
|
||||
JSON still has the ordinary object/array schema; synchronous callers materialize
|
||||
it explicitly. Bytes own their lifetime independently of the worker directory.
|
||||
"""
|
||||
from dataclasses import dataclass
|
||||
import json
|
||||
from pathlib import Path
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class NativeSeriesJson:
|
||||
data: bytes
|
||||
sample_count: int
|
||||
|
||||
def materialize(self) -> dict[str, list[float]]:
|
||||
return json.loads(self.data)
|
||||
|
||||
|
||||
def read_indexed_result(output: Path, index_path: Path) -> dict:
|
||||
index = json.loads(index_path.read_bytes())
|
||||
if not isinstance(index, dict) or type(index.get('version')) is not int or index['version'] != 1:
|
||||
raise ValueError('Unsupported native result index.')
|
||||
names = ('seriesStart', 'seriesEnd', 'resultBytes', 'sampleCount')
|
||||
if any(type(index.get(key)) is not int for key in names):
|
||||
raise ValueError('Invalid native result index integers.')
|
||||
start, end, length, count = (index[key] for key in names)
|
||||
if not (0 < start < end < length and count >= 0):
|
||||
raise ValueError('Invalid native result index bounds.')
|
||||
if output.stat().st_size != length:
|
||||
raise ValueError('Incomplete native result file.')
|
||||
raw = output.read_bytes()
|
||||
if (len(raw) != length or not raw[:start].endswith(b'"series":')
|
||||
or raw[start:start+1] != b'{' or raw[end-1:end] != b'}'
|
||||
or not raw[end:].startswith(b',"final":')):
|
||||
raise ValueError('Native result index does not match the output layout.')
|
||||
# The C writer supplies exact boundaries. No search through strings or numeric
|
||||
# arrays, and no large json.loads call: only diagnostics/final values are read.
|
||||
metadata = json.loads(raw[:start] + b'{}' + raw[end:])
|
||||
if not isinstance(metadata, dict) or metadata.get('series') != {}:
|
||||
raise ValueError('Invalid native result metadata.')
|
||||
metadata['series'] = NativeSeriesJson(raw[start:end], count)
|
||||
return metadata
|
||||
|
||||
|
||||
def serialize_result_parts(payload: dict) -> tuple[bytes, ...]:
|
||||
"""Return one JSON object in three byte segments, without copying its series.
|
||||
|
||||
The only raw position is payload.result.series, produced by our C writer;
|
||||
all model-provided labels, messages and keys use the standard JSON encoder.
|
||||
"""
|
||||
result = payload.get('result')
|
||||
series = result.get('series') if isinstance(result, dict) else None
|
||||
if not isinstance(series, NativeSeriesJson):
|
||||
return (json.dumps(payload, ensure_ascii=False, separators=(',', ':')).encode('utf-8'),)
|
||||
metadata = {key: value for key, value in result.items() if key != 'series'}
|
||||
outer = {key: value for key, value in payload.items() if key != 'result'}
|
||||
encoded = json.dumps(metadata, ensure_ascii=False, separators=(',', ':')).encode('utf-8')
|
||||
remainder = json.dumps(outer, ensure_ascii=False, separators=(',', ':')).encode('utf-8')
|
||||
prefix = b'{"result":' + encoded[:-1] + (b',' if metadata else b'') + b'"series":'
|
||||
suffix = b'}' + (b',' + remainder[1:] if outer else b'}')
|
||||
return prefix, series.data, suffix
|
||||
@@ -2,7 +2,10 @@
|
||||
from __future__ import annotations
|
||||
|
||||
from dataclasses import dataclass
|
||||
from typing import Literal
|
||||
from typing import TYPE_CHECKING, Literal
|
||||
|
||||
if TYPE_CHECKING:
|
||||
from app.simulation.native_codegen.transport import NativeSeriesJson
|
||||
from app.simulation.core.metadata import ResultVariableMetadata
|
||||
|
||||
SimulationRunStatus = Literal["completed", "cancelled", "failed"]
|
||||
@@ -30,11 +33,15 @@ class GenericSimulationResult:
|
||||
simulated_until: float
|
||||
requested_stop_time: float
|
||||
variables: tuple[ResultVariableMetadata, ...]
|
||||
series: dict[str, list[float]]
|
||||
series: dict[str, list[float]] | NativeSeriesJson
|
||||
final: dict[str, float]
|
||||
diagnostics: dict[str, object]
|
||||
|
||||
def as_dict(self) -> dict[str, object]:
|
||||
def as_dict(self, *, raw_series: bool = False) -> dict[str, object]:
|
||||
from app.simulation.native_codegen.transport import NativeSeriesJson
|
||||
series = self.series
|
||||
if isinstance(series, NativeSeriesJson) and not raw_series:
|
||||
series = series.materialize()
|
||||
return {
|
||||
"success": self.success,
|
||||
"status": self.status,
|
||||
@@ -43,7 +50,7 @@ class GenericSimulationResult:
|
||||
"simulatedUntil": self.simulated_until,
|
||||
"requestedStopTime": self.requested_stop_time,
|
||||
"variables": [variable.as_dict() for variable in self.variables],
|
||||
"series": self.series,
|
||||
"series": series,
|
||||
"final": self.final,
|
||||
"diagnostics": self.diagnostics,
|
||||
}
|
||||
Reference in new issue
Block a user