Files
SystemSimulationApp/docs/other/backups/2026-09-12-result-transfer/optimization.patch
T

1247 lines
74 KiB
Diff

--- a/app/main.py
+++ b/app/main.py
@@ -22,6 +22,7 @@
from fastapi.responses import FileResponse, HTMLResponse, StreamingResponse
from pydantic import BaseModel, ConfigDict, Field, ValidationError
+from app.result_compression import SimulationResultCompressionMiddleware
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
@@ -58,6 +59,7 @@
title="System Simulation ReactFlow App",
lifespan=_app_lifespan,
)
+app.add_middleware(SimulationResultCompressionMiddleware)
FRONTEND_DIST_DIR = Path(__file__).resolve().parent.parent / "frontend" / "dist"
PROJECT_STORAGE_DIR = Path(__file__).parent / "data" / "reactflow-projects"
SYSTEM_XML_SCHEMA_VERSION = "3"
@@ -73,7 +75,7 @@
"postprocessing": "正在整理采样结果",
"cancelled": "正在整理已终止仿真的部分结果",
"failed": "正在整理异常终止前的部分结果",
- "complete": "正在汇总仿真结果",
+ "complete": "计算已完成,正在接收仿真结果",
}
SIMULATION_STREAM_HEARTBEAT_SECONDS = 5.0
SIMULATION_TASK_RETENTION_SECONDS = 600.0
--- a/frontend/src/App.tsx
+++ b/frontend/src/App.tsx
@@ -11341,6 +11341,8 @@
throw new SimulationStreamError("浏览器未收到仿真进度数据流");
}
+ const contentEncoding = response.headers.get("Content-Encoding")?.trim().toLowerCase();
+ const requiresCompleteBody = Boolean(contentEncoding && contentEncoding !== "identity");
const reader = response.body.getReader();
let result: SimulationResult | null = null;
let activityWatchdog = createSimulationActivityWatchdog();
@@ -11411,6 +11413,10 @@
while (true) {
const { done, value } = await readSimulationStreamChunk(reader);
if (value) lineDecoder.write(value);
+ // Identity responses can finish at the terminal record. Encoded bodies
+ // must reach EOF so the browser validates gzip's trailer (or equivalent
+ // encoding integrity); decoded JSON can arrive before that validation.
+ if (result && !requiresCompleteBody) break;
if (done) {
lineDecoder.finish();
break;
@@ -11419,7 +11425,9 @@
} finally {
abortController.abort();
try {
- await reader.cancel();
+ // cancel() closes the reader locally before its underlying-source promise
+ // settles; remote cleanup must not delay publishing a completed result.
+ void reader.cancel().catch(() => {});
} catch {
// The stream may already be closed or aborted.
}
--- a/tests/manual/browser_stage_profile.mjs
+++ b/tests/manual/browser_stage_profile.mjs
@@ -1,4 +1,5 @@
-// Real production-page profiling. No route mocks, response cloning, or duplicate body parsing.
+// Real-page profiling; production by default, development only with explicit opt-in.
+// No route mocks, response cloning, or duplicate body parsing.
// All stage timestamps use the active document's performance.now(). A reload starts a new axis.
import { chromium } from '../../frontend/node_modules/playwright/index.mjs';
import fs from 'node:fs/promises';
@@ -6,11 +7,12 @@
import assert from 'node:assert/strict';
import { createHash } from 'node:crypto';
-const usage = `node tests/manual/browser_stage_profile.mjs --output DIR [--input tests/data/test-mql-8-corrected.json] [--url http://127.0.0.1:8011] [--runs 3 (0 for one smoke run)] [--mode both|profiled|control] [--deep] [--cpu-interval-us 1000] [--source-map-dir DIR] [--check]
+const usage = `node tests/manual/browser_stage_profile.mjs --output DIR [--input tests/data/test-mql-8-corrected.json] [--url http://127.0.0.1:8011] [--runs 3 (0 for one smoke run)] [--mode both|profiled|control] [--deep] [--cpu-interval-us 1000] [--source-map-dir DIR] [--allow-development] [--check]
Offline only: node tests/manual/browser_stage_profile.mjs --summarize-cpu-only EXISTING_DIRECTORY [--source-map-dir DIR] [--check]
Each mode runs one warmup followed by RUNS measured runs, sequentially. --check validates inputs without launching a browser.
Optional --deep (alias --cpu-profile) records renderer-main-thread .cpuprofile diagnostics separately from ordinary endpoint timing.
--source-map-dir accepts an offline hidden-source-map build only when its generated JS bytes exactly match served assets.
+--allow-development permits Vite development assets (for example forwarded port 5173); reports are labelled development and must not be mixed with production timings.
Run with the repository Node 24 and Chromium runtime library environment. No application code is modified.`;
const args = process.argv.slice(2);
if (args.includes('--help')) { console.log(usage); process.exit(0); }
@@ -18,6 +20,7 @@
for (let i = 0; i < args.length; i++) {
if (['--deep', '--cpu-profile'].includes(args[i])) { options.deep = true; continue; }
if (args[i] === '--check') { options.check = true; continue; }
+ if (args[i] === '--allow-development') { options.allowDevelopment = true; continue; }
const cliKey = args[i].replace(/^--/, '');
const key = ({ 'cpu-interval-us': 'cpuIntervalUs', 'source-map-dir': 'sourceMapDir', 'summarize-cpu-only': 'summarizeCpuOnly' })[cliKey] ?? cliKey;
if (!['input', 'output', 'url', 'runs', 'mode', 'cpuIntervalUs', 'sourceMapDir', 'summarizeCpuOnly'].includes(key) || !args[i + 1]) throw new Error(usage);
@@ -63,6 +66,8 @@
Object.assign(data, { marks: {}, requests: [], reads: [], parses: [], decodes: [], transactions: [], workers: [], runFailure: undefined,
downloads: [], longTasks: [], streamActive: false, activeRun: true });
performance.clearMarks();
+ // Resource entries for this run must not be displaced by Vite startup imports.
+ performance.clearResourceTimings();
};
data.armClick = (name, selector) => {
const listener = event => {
@@ -83,7 +88,10 @@
const failure = document.querySelector('.simulation-console-dock-progress.error, .simulation-console-progress.error');
if (failure) { data.runFailure = failure.textContent; mark('runFailure'); observer.disconnect(); return; }
const button = document.querySelector('button[aria-label="运行仿真"]');
- observedBusy ||= Boolean(button?.disabled || document.querySelector('.simulation-console-dock-progress.running, .simulation-console-progress.running'));
+ const runningProgress = document.querySelector('.simulation-console-dock-progress.running, .simulation-console-progress.running');
+ observedBusy ||= Boolean(button?.disabled || runningProgress);
+ const progressText = runningProgress?.textContent ?? '';
+ if (progressText.includes('正在汇总仿真结果') || progressText.includes('计算已完成,正在接收仿真结果')) mark('completeProgressDom');
const success = document.querySelector('.simulation-console-dock-progress.success, .simulation-console-progress.success');
if (observedBusy && button && !button.disabled && success?.textContent.includes('仿真完成')) {
mark('resultReadyDom');
@@ -151,7 +159,8 @@
.filter(e => /simulate-stream|simulation-results\/csv/.test(e.name))
.map(e => ({ name: e.name, startTime: e.startTime, requestStart: e.requestStart, responseStart: e.responseStart,
responseEnd: e.responseEnd, duration: e.duration, transferSize: e.transferSize,
- encodedBodySize: e.encodedBodySize, decodedBodySize: e.decodedBodySize })) });
+ encodedBodySize: e.encodedBodySize, decodedBodySize: e.decodedBodySize,
+ nextHopProtocol: e.nextHopProtocol, deliveryType: e.deliveryType ?? null })) });
if (location.hash === '#/results' && sessionStorage.getItem(storageKey)) {
data.watchDom('restoredResults', '.results-shell .results-system-panel');
}
@@ -169,6 +178,8 @@
if (kind === 'simulation') { mark('fetchStart', row.fetchStart); data.streamActive = true; }
return Reflect.apply(originalFetch, this, args).then(response => {
row.headers = performance.now(); row.status = response.status;
+ row.contentEncoding = response.headers.get('content-encoding');
+ row.contentLength = response.headers.get('content-length');
if (kind === 'simulation') mark('headers', row.headers);
if (response.body) bodies.set(response.body, row);
return response;
@@ -200,9 +211,14 @@
const parsed = Reflect.apply(nativeParse, this, args);
const end = performance.now();
if (parsed && ['progress', 'result', 'error'].includes(parsed.event)) {
- data.parses.push({ event: parsed.event, phase: parsed.phase, start, end, characters: typeof args[0] === 'string' ? args[0].length : null });
+ data.parses.push({ event: parsed.event, phase: parsed.phase, heartbeat: Boolean(parsed.heartbeat), start, end, characters: typeof args[0] === 'string' ? args[0].length : null });
+ if (parsed.event === 'progress' && parsed.phase === 'complete' && !parsed.heartbeat) {
+ mark('completeProgressParsed', end);
+ }
if (parsed.event === 'result') {
mark('resultParseStart', start); mark('resultParseEnd', end);
+ // A later suffix/newline/read may otherwise overwrite lastChunk.
+ if (data.marks.lastChunk !== undefined) mark('resultBoundaryReadAtParseStart', data.marks.lastChunk);
// This microtask is only a checkpoint after the current consumer continuation,
// not a claim that all EOF/finally/publish/React work has completed.
queueMicrotask(() => mark('resultConsumerMicrotaskCheckpoint'));
@@ -676,16 +692,17 @@
if (options.check) {
checkCpuSummaryContract();
+ checkFinalizationMeasurementContract();
if (options.summarizeCpuOnly) {
const manifest = JSON.parse(await fs.readFile(path.join(options.summarizeCpuOnly, 'summary.json')));
assert.ok(manifest.cpuDiagnostics?.profiles?.length);
console.log(JSON.stringify({ directory: options.summarizeCpuOnly, profiles: manifest.cpuDiagnostics.profiles.length,
- cpuClassificationContractPassed: true, browserLaunched: false }));
+ cpuClassificationContractPassed: true, finalizationMeasurementContractPassed: true, browserLaunched: false }));
} else {
console.log(JSON.stringify({ input: options.input, inputSha256: sha(inputText), nodes: project.nodes.length,
edges: project.edges.length, curveNodeId, mode: options.mode, deep: Boolean(options.deep), cpuIntervalUs: Number(options.cpuIntervalUs),
- sourceMapDir: options.sourceMapDir ?? null, warmupsPerMode: 1, measuredRunsPerMode: Number(options.runs),
- cpuClassificationContractPassed: true, browserLaunched: false }));
+ sourceMapDir: options.sourceMapDir ?? null, allowDevelopment: Boolean(options.allowDevelopment), warmupsPerMode: 1, measuredRunsPerMode: Number(options.runs),
+ cpuClassificationContractPassed: true, finalizationMeasurementContractPassed: true, browserLaunched: false }));
}
process.exit(0);
}
@@ -701,6 +718,9 @@
headersAndReads: 'fetch resolution and consumer read delivery. Outstanding-read intervals include backend production, transport and browser scheduling; not pure network time.',
unobservedCpu: 'NDJSON fragment scanning/join/trim and React handler CPU are not isolated. Their residual intervals can also include scheduling and cannot be attributed wholesale to parsing, transport or drawing.',
jsonParse: 'Only the original synchronous JSON.parse call, called once per application parse. No duplicate body read, decode, scan or parse.',
+ finalization: 'completeProgressDom is the first target text in the active running-progress DOM, observed in control and profiled modes; not a paint timestamp. completeProgressParsed is the first non-heartbeat complete progress JSON parse end, profiled only. The two origins are independent and may be absent or differently ordered. Endpoints include result parse/EOF/ready; null and negative deltas are retained. If the application finishes on the result event and cancels before EOF, EOF stays absent rather than being invented.',
+ responseEncoding: 'Content-Encoding/Content-Length are captured from real response headers without reading the body. Resource Timing encodedBodySize is compressed HTTP body size; decodedBodySize and reader byte counts are after content decoding. transferSize also includes browser-reported response-header overhead; none is a precise TLS/port-forward wire-byte measurement. Zero sizes may mean caching or unavailable/incomplete timing and are not replaced by decoded bytes. An early stream cancel may leave Resource Timing incomplete or absent.',
+ frontendEnvironment: 'Production assets are required by default. --allow-development permits observed Vite assets and labels the group development; development costs must not be combined with production groups. Asset hashes cover document tags and observed workers, not the entire Vite module graph.',
resultReady: 'First DOM observation of successful completion plus an enabled run button after busy state. React state/handler boundaries are not directly instrumented.',
indexedDb: 'Profiled: session pointer publication immediately after all save transactions commit. Control: pointer polling, up to 16 ms plus scheduling delay. Transaction windows also include asynchronous waiting and may include old-cache cleanup.',
render: 'First visible DOM, then two requestAnimationFrame callbacks (paint opportunity, not GPU completion); stable means scoped DOM quiet for 120 ms followed by two frames.',
@@ -716,6 +736,7 @@
const errors = [];
const browser = await chromium.launch({ headless: true });
const evidence = { input: path.resolve(options.input), inputSha256: sha(inputText), baseURL: options.url,
+ allowDevelopment: Boolean(options.allowDevelopment), frontendEnvironment: null,
browser: browser.version(), node: process.version, deep: Boolean(options.deep), cpuIntervalUs: options.deep ? Number(options.cpuIntervalUs) : null, initialNavigation: {}, scriptSha256: sha(await fs.readFile(new URL(import.meta.url))), definitions, rows, errors, servedAssets: [] };
const writeSummary = () => fs.writeFile(path.join(options.output, 'summary.json'), JSON.stringify(evidence, null, 2));
const waitMark = async (page, name) => {
@@ -728,16 +749,40 @@
const watch = (page, name, selector) => page.evaluate(({ name, selector }) => window.__stageProfile.watchDom(name, selector), { name, selector });
const clickTab = async (page, name) => { await page.getByRole('tab', { name }).click(); };
const resultDigest = result => sha(JSON.stringify(result));
-const delta = (m, a, b) => m[a] === undefined || m[b] === undefined ? null : m[b] - m[a];
+function delta(m, a, b) { return m[a] === undefined || m[b] === undefined ? null : m[b] - m[a]; }
function metrics(trace) {
const m = trace.marks;
const csvAnchor = trace.downloads.find(d => d.name.endsWith('.csv'))?.anchorClick;
const resultAnchor = trace.downloads.find(d => d.name.endsWith('.simresult'))?.anchorClick;
const csvRequest = trace.requests.find(r => r.kind === 'csv');
+ const simulationRequest = trace.requests.find(r => r.kind === 'simulation');
+ const simulationResponse = trace.simulationResponse ?? simulationRequest;
+ const resources = (trace.resources ?? []).filter(r => r.name.includes('/api/system-xml/simulate-stream'));
+ const resource = resources.length === 1 ? resources[0] : null;
+ const contentLength = simulationResponse?.contentLength;
+ const finalization = {};
+ for (const [label, start] of [['completeDom', 'completeProgressDom'], ['completeParsed', 'completeProgressParsed']]) {
+ for (const [endpoint, end] of [['ReadyDom', 'resultReadyDom'], ['ReadyPaintOpportunity', 'resultReadyPaintOpportunity'],
+ ['ParseStart', 'resultParseStart'], ['ParseEnd', 'resultParseEnd'], ['Eof', 'streamEof'],
+ ['ResultBoundaryRead', 'resultBoundaryReadAtParseStart']]) {
+ finalization[`${label}To${endpoint}Ms`] = delta(m, start, end);
+ }
+ }
const worker = trace.workers?.[0];
const workerStart = worker?.posts.find(p => p.type === 'start');
const workerFinish = worker?.posts.find(p => p.type === 'finish');
return {
+ ...finalization,
+ completeParsedToDomMs: delta(m, 'completeProgressParsed', 'completeProgressDom'),
+ resultBoundaryReadToParseStartMs: delta(m, 'resultBoundaryReadAtParseStart', 'resultParseStart'),
+ responseContentEncoding: simulationResponse?.contentEncoding ?? null,
+ responseContentLengthBytes: typeof contentLength === 'string' && /^\d+$/.test(contentLength) ? Number(contentLength) : null,
+ responseTransferEncoding: simulationResponse?.transferEncoding ?? null,
+ resourceTimingSimulationEntries: resources.length,
+ resourceTransferBytes: resource?.transferSize ?? null,
+ resourceEncodedBodyBytes: resource?.encodedBodySize ?? null,
+ resourceDecodedBodyBytes: resource?.decodedBodySize ?? null,
+ resourceDecodedToEncodedRatio: resource?.encodedBodySize > 0 && resource?.decodedBodySize > 0 ? resource.decodedBodySize / resource.encodedBodySize : null,
importToDomMs: delta(m, 'importChange', 'importReadyDom'),
importToPaintOpportunityMs: delta(m, 'importChange', 'importReadyPaintOpportunity'),
clickToFetchMs: delta(m, 'runClick', 'fetchStart'), fetchToHeadersMs: delta(m, 'fetchStart', 'headers'),
@@ -773,6 +818,49 @@
synchronousStreamDecodeMs: trace.profiled ? trace.decodes.reduce((n, r) => n + r.end - r.start, 0) : null,
};
}
+function classifyFrontendEnvironment(assetUrls, allowDevelopment) {
+ const environment = assetUrls.some(url => /@vite\/client|@react-refresh|\/src\//.test(url)) ? 'development' : 'production';
+ if (environment === 'development' && !allowDevelopment) {
+ throw new Error('Expected a production build, found Vite development assets. Use --allow-development for a separately labelled development measurement.');
+ }
+ return environment;
+}
+
+function checkFinalizationMeasurementContract() {
+ const trace = { profiled: true, marks: { completeProgressParsed: 8, completeProgressDom: 10,
+ resultBoundaryReadAtParseStart: 35, resultParseStart: 40, resultParseEnd: 100, resultReadyDom: 120 },
+ requests: [], reads: [], parses: [], decodes: [], workers: [], downloads: [],
+ simulationResponse: { contentEncoding: 'gzip', contentLength: '100' },
+ resources: [{ name: 'http://localhost/api/system-xml/simulate-stream', transferSize: 120, encodedBodySize: 100, decodedBodySize: 1000 }] };
+ const row = metrics(trace);
+ assert.equal(row.completeDomToReadyDomMs, 110);
+ assert.equal(row.completeParsedToReadyDomMs, 112);
+ assert.equal(row.completeParsedToDomMs, 2);
+ assert.equal(row.completeParsedToParseStartMs, 32);
+ assert.equal(row.completeParsedToParseEndMs, 92);
+ assert.equal(row.completeParsedToEofMs, null, 'Early result publication must not invent EOF.');
+ assert.equal(row.resultBoundaryReadToParseStartMs, 5);
+ assert.equal(row.responseContentEncoding, 'gzip');
+ assert.equal(row.responseContentLengthBytes, 100);
+ assert.equal(row.resourceEncodedBodyBytes, 100);
+ assert.equal(row.resourceDecodedBodyBytes, 1000);
+ assert.equal(row.resourceDecodedToEncodedRatio, 10);
+ assert.equal(metrics({ ...trace, marks: { ...trace.marks, completeProgressDom: 60 } }).completeDomToParseStartMs, -20);
+ const control = metrics({ ...trace, profiled: false, marks: { completeProgressDom: 10, resultReadyDom: 120 }, resources: [] });
+ assert.equal(control.completeDomToReadyDomMs, 110);
+ assert.equal(control.completeParsedToReadyDomMs, null);
+ assert.equal(control.completeDomToParseStartMs, null);
+ assert.equal(control.streamBytes, null);
+ assert.equal(control.resourceEncodedBodyBytes, null);
+ const zero = metrics({ ...trace, resources: [{ ...trace.resources[0], transferSize: 0, encodedBodySize: 0 }] });
+ assert.equal(zero.resourceEncodedBodyBytes, 0);
+ assert.equal(zero.resourceDecodedToEncodedRatio, null, 'Missing wire-size evidence must not be replaced by decoded size.');
+ assert.equal(metrics({ ...trace, resources: [...trace.resources, ...trace.resources] }).resourceEncodedBodyBytes, null);
+ assert.equal(classifyFrontendEnvironment(['http://localhost/assets/index-abc.js'], false), 'production');
+ assert.equal(classifyFrontendEnvironment(['http://localhost/@vite/client'], true), 'development');
+ assert.throws(() => classifyFrontendEnvironment(['http://localhost/src/main.tsx'], false), /--allow-development/);
+}
+
try {
for (const mode of options.mode === 'both' ? ['control', 'profiled'] : [options.mode]) {
const context = await browser.newContext({ viewport: { width: 1600, height: 1000 }, acceptDownloads: true });
@@ -781,12 +869,25 @@
const cpuRecorder = options.deep ? await createCpuRecorder(context, page, Number(options.cpuIntervalUs)) : null;
page.setDefaultTimeout(30000);
const simulationRequests = [];
+ const requestRecords = new WeakMap();
const workerUrls = new Set();
page.on('worker', worker => workerUrls.add(worker.url()));
page.on('request', request => {
if (request.url().includes('/api/system-xml/simulate-stream')) {
- simulationRequests.push({ url: request.url(), simulationId: request.headers()['x-simulation-id'] ?? null });
+ const row = { url: request.url(), simulationId: request.headers()['x-simulation-id'] ?? null };
+ simulationRequests.push(row);
+ requestRecords.set(request, row);
}
+ });
+ // Passive header metadata in both modes; never response.body/text/json.
+ page.on('response', response => {
+ const row = requestRecords.get(response.request());
+ if (!row) return;
+ const headers = response.headers();
+ row.response = { status: response.status(), contentEncoding: headers['content-encoding'] ?? null,
+ contentLength: headers['content-length'] ?? null, transferEncoding: headers['transfer-encoding'] ?? null,
+ contentType: headers['content-type'] ?? null, vary: headers.vary ?? null,
+ fromServiceWorker: response.fromServiceWorker(), source: 'playwright-response-headers' };
});
page.on('pageerror', error => errors.push({ mode, error: String(error) }));
page.on('dialog', dialog => { errors.push({ mode, dialog: dialog.message() }); void dialog.dismiss(); });
@@ -797,16 +898,18 @@
appControlObservedAt: performance.now(), navigation: performance.getEntriesByType('navigation').map(e => e.toJSON()) }));
const assetUrls = await page.evaluate(() => [...document.querySelectorAll('script[src],link[rel="stylesheet"][href]')]
.map(e => e.src || e.href));
- if (assetUrls.some(url => /@vite\/client|\/src\//.test(url))) throw new Error('Expected a production build, found Vite development assets.');
+ const frontendEnvironment = classifyFrontendEnvironment(assetUrls, Boolean(options.allowDevelopment));
+ if (evidence.frontendEnvironment !== null) assert.equal(evidence.frontendEnvironment, frontendEnvironment, 'Frontend environment changed between modes.');
+ evidence.frontendEnvironment = frontendEnvironment;
for (const url of assetUrls) {
const response = await context.request.get(url);
assert.ok(response.ok(), `Asset HTTP ${response.status()}: ${url}`);
const bytes = await response.body();
const previous = evidence.servedAssets.find(asset => asset.url === url);
- if (previous) assert.equal(previous.sha256, sha(bytes), 'Production asset changed between modes.');
+ if (previous) assert.equal(previous.sha256, sha(bytes), 'Frontend asset changed between modes.');
else evidence.servedAssets.push({ url, sha256: sha(bytes), bytes: bytes.length });
}
- assert.ok(evidence.servedAssets.length, 'No production assets found.');
+ assert.ok(evidence.servedAssets.length, 'No frontend assets found.');
evidence.buildAssetSetSha256 = sha(JSON.stringify(evidence.servedAssets
.map(({ url, ...asset }) => ({ path: new URL(url).pathname, ...asset }))
.sort((a, b) => a.path.localeCompare(b.path))));
@@ -871,6 +974,11 @@
await saveDownload('下载结果 CSV', 'csvClick', 'csvDownloadSaved', 'result.csv');
await saveDownload('下载结果文件', 'resultFileClick', 'resultFileDownloadSaved', 'result.simresult');
const trace = await page.evaluate(() => window.__stageProfile.snapshot());
+ const requests = simulationRequests.slice(requestOffset);
+ assert.equal(requests.length, 1, 'Expected exactly one real simulation request.');
+ assert.ok(requests[0].simulationId, 'Missing X-Simulation-Id correlation key.');
+ assert.ok(requests[0].response, 'Missing simulation response-header evidence.');
+ trace.simulationResponse = { simulationId: requests[0].simulationId, ...requests[0].response };
if (cpuRecorder) {
const capture = await cpuRecorder.stop();
const file = path.join(runDir, 'interaction.cpuprofile');
@@ -908,10 +1016,7 @@
const restoreTrace = await page.evaluate(() => ({ ...window.__stageProfile.snapshot(),
navigation: performance.getEntriesByType('navigation').map(e => e.toJSON()) }));
await fs.writeFile(path.join(runDir, 'restore-trace.json'), JSON.stringify(restoreTrace, null, 2));
- const requests = simulationRequests.slice(requestOffset);
- assert.equal(requests.length, 1, 'Expected exactly one real simulation request.');
- assert.ok(requests[0].simulationId, 'Missing X-Simulation-Id correlation key.');
- const row = { mode, deep: Boolean(options.deep), run, warmup: run === 0, simulationId: requests[0].simulationId, timeOrigin: trace.timeOrigin, ...metrics(trace),
+ const row = { mode, frontendEnvironment, deep: Boolean(options.deep), run, warmup: run === 0, simulationId: requests[0].simulationId, timeOrigin: trace.timeOrigin, ...metrics(trace),
restoreNavigationToDomMs: restoreTrace.marks.restoredResultsDom,
restoreNavigationToPaintOpportunityMs: restoreTrace.marks.restoredResultsPaintOpportunity,
restoreNavigationToDomStableMs: restoreTrace.marks.restoredResultsStable,
--- a/requirements.txt
+++ b/requirements.txt
@@ -4,3 +4,6 @@
lxml>=5,<7
pydantic>=2,<3
uvicorn[standard]
+
+# Streaming gzip must flush progress and offload large compression blocks.
+starlette>=1.6,<2
--- a/tests/test_native_result_transport.py
+++ b/tests/test_native_result_transport.py
@@ -1,6 +1,9 @@
"""The HTTP fast path must preserve real native values and task semantics."""
from dataclasses import replace
+import gzip
import json
+import struct
+import zlib
from pathlib import Path
import tempfile
from types import SimpleNamespace
@@ -23,7 +26,7 @@
"""Exercise real routing/response bodies without an optional HTTP client dependency."""
def __init__(self, application): self.application = application
def post(self, path, *, content=b'', headers=None): return self.request('POST', path, content, headers)
- def get(self, path): return self.request('GET', path, b'', None)
+ def get(self, path, *, headers=None): return self.request('GET', path, b'', headers)
def request(self, method, path, content, headers):
async def run():
messages = []
@@ -37,7 +40,11 @@
await self.application(scope,receive,send)
status = next(m['status'] for m in messages if m['type']=='http.response.start')
body = b''.join(m.get('body',b'') for m in messages if m['type']=='http.response.body')
- return SimpleNamespace(status_code=status,content=body,json=lambda:json.loads(body))
+ response_headers = {key.decode('latin-1'): value.decode('latin-1')
+ for key, value in next(m['headers'] for m in messages if m['type']=='http.response.start')}
+ chunks = tuple(m.get('body', b'') for m in messages if m['type']=='http.response.body')
+ return SimpleNamespace(status_code=status, content=body, headers=response_headers,
+ chunks=chunks, json=lambda:json.loads(body))
return asyncio.run(run())
@@ -95,6 +102,71 @@
for key in ('series', 'final', 'variables', 'model', 'simulation'):
self.assertEqual(synchronous[key], result[key])
+ def test_real_http_gzip_stream_and_retained_task_preserve_identity_results(self):
+ client = AsgiClient(app)
+ ident = 'transport-gzip-'+uuid4().hex
+ def load(data):
+ return json.loads(data, parse_int=lambda token: -0.0 if token == '-0' else int(token))
+ def decode(response, *, progress=False):
+ self.assertEqual(response.status_code, 200)
+ self.assertEqual(response.headers['content-encoding'], 'gzip')
+ self.assertIn('accept-encoding', response.headers['vary'].lower())
+ decoder = zlib.decompressobj(zlib.MAX_WBITS + 16)
+ parts = []
+ first_progress = False
+ for index, chunk in enumerate(response.chunks):
+ part = decoder.decompress(chunk)
+ parts.append(part)
+ if progress and not first_progress and b'\n' in part:
+ self.assertEqual(load(part.split(b'\n', 1)[0])['event'], 'progress')
+ self.assertLess(index, len(response.chunks)-1)
+ self.assertFalse(decoder.eof, 'Progress must precede gzip stream completion')
+ first_progress = True
+ parts.append(decoder.flush())
+ decoded = b''.join(parts)
+ self.assertTrue(decoder.eof, 'The real HTTP response needs its complete gzip trailer')
+ self.assertEqual(decoder.unused_data, b'')
+ self.assertEqual(decoder.unconsumed_tail, b'')
+ self.assertEqual(gzip.decompress(response.content), decoded)
+ if progress:
+ self.assertTrue(first_progress)
+ return decoded
+ compressed = client.post('/api/system-xml/simulate-stream', content=self.xml,
+ headers={'X-Simulation-Id': ident, 'Accept-Encoding': 'gzip'})
+ events = [load(line) for line in decode(compressed, progress=True).splitlines()]
+ results = [event['result'] for event in events if event['event'] == 'result']
+ self.assertEqual(len(results), 1)
+ result = results[0]
+ self.assertTrue(result['success'])
+ self.assertEqual(result['simulatedUntil'], .1)
+ identity = client.post('/api/system-xml/simulate-stream', content=self.xml,
+ headers={'X-Simulation-Id': ident+'-identity', 'Accept-Encoding': 'identity'})
+ self.assertNotIn('content-encoding', identity.headers)
+ identity_result = next(event['result'] for event in map(load, identity.content.splitlines())
+ if event['event'] == 'result')
+ for key in ('success', 'status', 'partial', 'simulatedUntil', 'requestedStopTime',
+ 'variables', 'model', 'simulation'):
+ self.assertEqual(result[key], identity_result[key], key)
+ self.assertEqual(result['diagnostics']['integration']['totals'],
+ identity_result['diagnostics']['integration']['totals'])
+ # Compare numeric bits, including signed zero, independently of gzip's
+ # byte oracle; timings legitimately vary between these two actual runs.
+ self.assertEqual(result['series'].keys(), identity_result['series'].keys())
+ for key, values in result['series'].items():
+ expected = identity_result['series'][key]
+ self.assertEqual(len(values), len(expected), key)
+ self.assertEqual(b''.join(struct.pack('!d', value) for value in values),
+ b''.join(struct.pack('!d', value) for value in expected), key)
+ self.assertEqual(result['final'].keys(), identity_result['final'].keys())
+ for key, value in result['final'].items():
+ self.assertEqual(struct.pack('!d', value), struct.pack('!d', identity_result['final'][key]), key)
+ retained = client.get('/api/system-xml/simulations/'+ident, headers={'Accept-Encoding': 'gzip'})
+ retained_bytes = decode(retained)
+ retained_identity = client.get('/api/system-xml/simulations/'+ident,
+ headers={'Accept-Encoding': 'identity'})
+ self.assertEqual(retained_bytes, retained_identity.content)
+ self.assertEqual(load(retained_bytes)['result'], result)
+
def test_cancelled_raw_stream_keeps_partial_result_and_public_status(self):
for reason, status in [('user', 'stopped'), ('stalled', 'stalled')]:
task = _register_simulation_task('transport-'+uuid4().hex)
--- /dev/null
+++ b/app/result_compression.py
@@ -0,0 +1,84 @@
+"""Lossless, promptly flushed compression for large simulation responses.
+
+Starlette >= 1.6 flushes each streaming body and moves large compression
+blocks off the event loop. Bound its input so transmission can overlap the
+encoding of the next block, including when the native series is one big body.
+"""
+
+from __future__ import annotations
+
+import re
+
+from starlette.datastructures import Headers
+from starlette.middleware.gzip import GZipResponder, IdentityResponder
+from starlette.types import ASGIApp, Message, Receive, Scope, Send
+
+
+RESULT_COMPRESSION_CHUNK_BYTES = 256 * 1024
+_QUALITY = re.compile(r"(?:0(?:\.\d{0,3})?|1(?:\.0{0,3})?)\Z")
+
+
+def _accepts_gzip(headers: Headers) -> bool:
+ qualities: dict[str, float] = {}
+ for header in headers.getlist("accept-encoding"):
+ for item in header.split(","):
+ coding, *parameters = item.lower().strip().split(";")
+ coding = coding.strip()
+ quality = 1.0
+ for parameter in parameters:
+ key, separator, value = parameter.strip().partition("=")
+ if key.strip() == "q":
+ value = value.strip()
+ quality = min(quality, float(value) if separator and _QUALITY.fullmatch(value) else 0.0)
+ # A repeated prohibition takes precedence over another entry.
+ qualities[coding] = min(qualities.get(coding, 1.0), quality)
+ gzip_quality = qualities.get("gzip", qualities.get("*", 0.0))
+ return gzip_quality > 0.0 and gzip_quality >= qualities.get("identity", 0.0)
+
+
+def _is_result_request(scope: Scope) -> bool:
+ path = scope.get("path", "")
+ method = scope.get("method", "")
+ if method == "POST":
+ return path in {"/api/system-xml/simulate", "/api/system-xml/simulate-stream"}
+ prefix = "/api/system-xml/simulations/"
+ if method == "GET" and path.startswith(prefix):
+ identifier = path[len(prefix):]
+ return bool(identifier) and "/" not in identifier
+ return False
+
+
+class SimulationResultCompressionMiddleware:
+ def __init__(self, app: ASGIApp) -> None:
+ self.app = app
+
+ async def __call__(self, scope: Scope, receive: Receive, send: Send) -> None:
+ if scope["type"] != "http" or not _is_result_request(scope):
+ await self.app(scope, receive, send)
+ return
+ if not _accepts_gzip(Headers(scope=scope)):
+ await IdentityResponder(self.app, minimum_size=500)(scope, receive, send)
+ return
+
+ async def bounded_app(scope: Scope, receive: Receive, send: Send) -> None:
+ can_split = True
+
+ async def send_bounded(message: Message) -> None:
+ nonlocal can_split
+ if message["type"] == "http.response.start":
+ can_split = message["status"] != 206 and "content-encoding" not in Headers(raw=message["headers"])
+ body = message.get("body", b"")
+ if message["type"] == "http.response.body" and can_split and len(body) > RESULT_COMPRESSION_CHUNK_BYTES:
+ for offset in range(0, len(body), RESULT_COMPRESSION_CHUNK_BYTES):
+ end = offset + RESULT_COMPRESSION_CHUNK_BYTES
+ await send({
+ **message,
+ "body": body[offset:end],
+ "more_body": end < len(body) or message.get("more_body", False),
+ })
+ else:
+ await send(message)
+
+ await self.app(scope, receive, send_bounded)
+
+ await GZipResponder(bounded_app, minimum_size=500, compresslevel=1)(scope, receive, send)
--- /dev/null
+++ b/frontend/tests/e2e/simulation-stream-terminal.spec.ts
@@ -0,0 +1,183 @@
+import { expect, test, type Page } from "@playwright/test";
+import { expandSimulationConsole, prepareApp, resultSnapshot, wideProject } from "./fixtures";
+
+async function prepareOpenStream(page: Page, cancelBehavior: "pending" | "reject" = "pending", encoding?: string) {
+ await prepareApp(page);
+ await page.addInitScript(({ cancelBehavior, encoding }) => {
+ const nativeFetch = window.fetch.bind(window);
+ const probe = {
+ body: null as ReadableStream<Uint8Array> | null,
+ signal: null as AbortSignal | null,
+ enqueue: (_text: string) => {},
+ close: () => {},
+ fail: (_message: string) => {},
+ closedByServer: false,
+ cancellations: 0,
+ unhandled: [] as string[],
+ };
+ (window as any).__terminalStream = probe;
+ window.addEventListener("unhandledrejection", (event) => probe.unhandled.push(String(event.reason)));
+ window.fetch = async (input, init) => {
+ if (!String(input).includes("/api/system-xml/simulate-stream")) return nativeFetch(input, init);
+ const encoder = new TextEncoder();
+ probe.signal = init?.signal ?? null;
+ probe.body = new ReadableStream<Uint8Array>({
+ start(controller) {
+ probe.enqueue = (text) => controller.enqueue(encoder.encode(text));
+ probe.close = () => { probe.closedByServer = true; controller.close(); };
+ probe.fail = (message) => controller.error(new TypeError(message));
+ },
+ cancel() {
+ probe.cancellations++;
+ // The response does not reach EOF; underlying cleanup may also never
+ // settle or reject. Neither is allowed to block a terminal result.
+ return cancelBehavior === "pending"
+ ? new Promise<void>(() => {})
+ : Promise.reject(new Error("test transport cancellation rejection"));
+ },
+ });
+ const headers = new Headers({ "Content-Type": "application/x-ndjson" });
+ if (encoding) headers.set("Content-Encoding", encoding);
+ return new Response(probe.body, { headers });
+ };
+ }, { cancelBehavior, encoding });
+ await page.goto("/");
+ await page.locator('input[type="file"]').setInputFiles({
+ name: "terminal-stream.json", mimeType: "application/json",
+ buffer: Buffer.from(JSON.stringify({ ...wideProject, edges: [...wideProject.edges, {
+ id: "close-test-loop", source: "generic_sensor_4", target: "generic_sensor_1",
+ sourceHandle: "port_b", targetHandle: "port_a", data: { isContactEdge: false },
+ }] })),
+ });
+ await expect(page.getByRole("textbox", { name: "工程", exact: true })).toHaveValue("terminal-stream");
+ await page.getByRole("button", { name: "运行仿真", exact: true }).click();
+ await expect.poll(() => page.evaluate(() => Boolean((window as any).__terminalStream.body))).toBe(true);
+ await expandSimulationConsole(page);
+}
+
+async function expectReleasedStream(page: Page) {
+ expect(await page.evaluate(() => {
+ const probe = (window as any).__terminalStream;
+ return { aborted: probe.signal.aborted, locked: probe.body.locked,
+ cancellations: probe.cancellations, closedByServer: probe.closedByServer, unhandled: probe.unhandled };
+ })).toEqual({ aborted: true, locked: false, cancellations: 1, closedByServer: false, unhandled: [] });
+}
+
+for (const status of ["completed", "stopped"] as const) {
+ test(`${status} result is published before EOF and transport cleanup, preserving every value`, async ({ page }) => {
+ const errors: string[] = [];
+ page.on("pageerror", (error) => errors.push(error.message));
+ await prepareOpenStream(page, status === "completed" ? "pending" : "reject", status === "stopped" ? "identity" : undefined);
+ const result = {
+ ...resultSnapshot.result, status, success: status === "completed", partial: status === "stopped",
+ simulatedUntil: status === "completed" ? 10 : 4,
+ series: { time: status === "completed" ? [0, 5, 10] : [0, 2, 4],
+ "generic_sensor_1.value": [Number.MIN_VALUE, -1.2345678901234567, Number.MAX_VALUE] },
+ final: { "generic_sensor_1.value": Number.MAX_VALUE },
+ };
+ const panel = page.getByRole("complementary", { name: "仿真控制台", exact: true });
+ const run = page.getByRole("button", { name: "运行仿真", exact: true });
+ // The complete JSON text alone is not yet a complete NDJSON record.
+ await page.evaluate((result) => {
+ (window as any).__terminalStream.enqueue(JSON.stringify({ event: "progress", phase: "complete",
+ progress: 100, simulatedTime: result.simulatedUntil, totalTime: 10, message: "正在汇总仿真结果" }) + "\n" +
+ JSON.stringify({ event: "result", result }));
+ }, result);
+ await expect(panel).toContainText("正在汇总仿真结果");
+ await expect(run).toBeDisabled();
+ await page.evaluate(() => (window as any).__terminalStream.enqueue("\n"));
+ // This deadline is below the 30-second idle timeout; EOF is never sent.
+ await expect(run).toBeEnabled({ timeout: 3000 });
+ await expect(panel).toContainText(status === "completed" ? "仿真完成,已生成新的结果" : "已手动终止");
+ await expect(panel).not.toContainText("仿真失败");
+ await expectReleasedStream(page);
+ await page.getByRole("tab", { name: /^结果/ }).click();
+ const pendingDownload = page.waitForEvent("download");
+ await page.getByRole("button", { name: "下载结果文件", exact: true }).click();
+ const stream = await (await pendingDownload).createReadStream();
+ expect(stream).not.toBeNull();
+ const chunks: Buffer[] = [];
+ for await (const chunk of stream!) chunks.push(Buffer.from(chunk));
+ const exported = JSON.parse(Buffer.concat(chunks).toString("utf8"));
+ expect(exported.snapshot.result).toEqual(result);
+ expect(errors).toEqual([]);
+ await expectReleasedStream(page);
+ });
+}
+
+test("an error before any terminal result releases an open stream without waiting for cleanup", async ({ page }) => {
+ const errors: string[] = [];
+ page.on("pageerror", (error) => errors.push(error.message));
+ await prepareOpenStream(page);
+ await page.evaluate(() => (window as any).__terminalStream.enqueue(JSON.stringify({
+ event: "error", message: "测试结果生成失败", status: 422,
+ }) + "\n"));
+ await expect(page.getByRole("button", { name: "运行仿真", exact: true })).toBeEnabled({ timeout: 3000 });
+ await expect(page.getByRole("complementary", { name: "仿真控制台", exact: true })).toContainText("测试结果生成失败");
+ await expectReleasedStream(page);
+ expect(errors).toEqual([]);
+});
+
+for (const truncated of [false, true]) {
+ test(`EOF ${truncated ? "rejects truncated JSON" : "accepts a complete result without a trailing newline"}`, async ({ page }) => {
+ await prepareOpenStream(page);
+ await page.evaluate(({ result, truncated }) => {
+ const probe = (window as any).__terminalStream;
+ const line = JSON.stringify({ event: "result", result });
+ probe.enqueue(truncated ? line.slice(0, -1) : line);
+ probe.close();
+ }, { result: resultSnapshot.result, truncated });
+ await expect(page.getByRole("button", { name: "运行仿真", exact: true })).toBeEnabled({ timeout: 3000 });
+ const panel = page.getByRole("complementary", { name: "仿真控制台", exact: true });
+ await expect(panel).toContainText(truncated ? "仿真失败" : "仿真完成,已生成新的结果");
+ expect(await page.evaluate(() => {
+ const probe = (window as any).__terminalStream;
+ return { aborted: probe.signal.aborted, locked: probe.body.locked, closedByServer: probe.closedByServer };
+ })).toEqual({ aborted: true, locked: false, closedByServer: true });
+ });
+}
+
+for (const invalidTrailer of [false, true]) {
+ test(`encoded result ${invalidTrailer ? "rejects a later decoding error" : "waits for validated EOF before publication"}`, async ({ page }) => {
+ const errors: string[] = [];
+ page.on("pageerror", (error) => errors.push(error.message));
+ await prepareOpenStream(page, "pending", "gzip");
+ // Fetch exposes decoded bytes while retaining Content-Encoding. This mock
+ // represents the browser delivering JSON before validating gzip's trailer;
+ // gzip wire encoding itself belongs to the HTTP/backend integration tests.
+ await page.evaluate((result) => {
+ (window as any).__terminalStream.enqueue(JSON.stringify({ event: "progress", phase: "complete",
+ progress: 100, simulatedTime: 10, totalTime: 10, message: "正在汇总仿真结果" }) + "\n" +
+ JSON.stringify({ event: "result", result }) + "\n");
+ }, resultSnapshot.result);
+ const panel = page.getByRole("complementary", { name: "仿真控制台", exact: true });
+ const run = page.getByRole("button", { name: "运行仿真", exact: true });
+ await expect(panel).toContainText("正在汇总仿真结果");
+ await expect(run).toBeDisabled();
+ expect(await page.evaluate(() => {
+ const probe = (window as any).__terminalStream;
+ return { locked: probe.body.locked, aborted: probe.signal.aborted };
+ })).toEqual({ locked: true, aborted: false });
+ await page.evaluate((invalidTrailer) => {
+ const probe = (window as any).__terminalStream;
+ if (invalidTrailer) probe.fail("gzip trailer checksum failed");
+ else probe.close();
+ }, invalidTrailer);
+ await expect(run).toBeEnabled({ timeout: 3000 });
+ if (invalidTrailer) {
+ await expect(panel).toContainText("gzip trailer checksum failed");
+ await expect(panel).not.toContainText("仿真完成,已生成新的结果");
+ expect(await page.evaluate(() => sessionStorage.getItem("system-simulation-flow:latest-result"))).toBeNull();
+ } else {
+ await expect(panel).toContainText("仿真完成,已生成新的结果");
+ await expect(panel).not.toContainText("仿真失败");
+ await page.getByRole("tab", { name: /^结果/ }).click();
+ await expect(page.getByRole("button", { name: "下载结果文件", exact: true })).toBeVisible();
+ }
+ expect(await page.evaluate(() => {
+ const probe = (window as any).__terminalStream;
+ return { locked: probe.body.locked, aborted: probe.signal.aborted, unhandled: probe.unhandled };
+ })).toEqual({ locked: false, aborted: true, unhandled: [] });
+ expect(errors).toEqual([]);
+ });
+}
--- /dev/null
+++ b/tests/test_simulation_result_compression.py
@@ -0,0 +1,251 @@
+"""Lossless, live ASGI result compression without a model or HTTP client.
+
+The wire oracle is Python's independent gzip/zlib decoder. Tests assert emitted
+bytes and request/response semantics, not the compressor's internal algorithm.
+"""
+from __future__ import annotations
+
+import asyncio
+import gzip
+import hashlib
+import unittest
+import zlib
+
+from app.result_compression import SimulationResultCompressionMiddleware
+
+
+STREAM_PATH = '/api/system-xml/simulate-stream'
+CHUNK_BYTES = 256 * 1024
+
+
+def http_scope(method='POST', path=STREAM_PATH, accept_encoding='gzip'):
+ return {
+ 'type': 'http', 'asgi': {'version': '3.0', 'spec_version': '2.4'},
+ 'http_version': '1.1', 'method': method, 'scheme': 'http',
+ 'path': path, 'raw_path': path.encode(), 'query_string': b'',
+ 'root_path': '', 'server': ('testserver', 80),
+ 'client': ('127.0.0.1', 1234),
+ 'headers': ([] if accept_encoding is None else
+ [(b'accept-encoding', value.encode('ascii'))
+ for value in ([accept_encoding] if isinstance(accept_encoding, str) else accept_encoding)]),
+ }
+
+
+def copied_message(message):
+ result = dict(message)
+ if 'headers' in result:
+ result['headers'] = list(result['headers'])
+ return result
+
+
+def response_headers(messages):
+ return dict(next(m['headers'] for m in messages if m['type'] == 'http.response.start'))
+
+
+def response_bodies(messages):
+ return [m for m in messages if m['type'] == 'http.response.body']
+
+
+async def receive():
+ return {'type': 'http.request', 'body': b'', 'more_body': False}
+
+
+class SimulationResultCompressionTests(unittest.IsolatedAsyncioTestCase):
+ async def request(self, body, *, method='POST', path=STREAM_PATH,
+ accept_encoding='gzip', status=200, headers=None):
+ original = [
+ {'type': 'http.response.start', 'status': status,
+ 'headers': list(headers if headers is not None else [
+ (b'content-type', b'application/json'),
+ (b'content-length', str(len(body)).encode()),
+ ])},
+ {'type': 'http.response.body', 'body': body, 'more_body': False},
+ ]
+ async def application(scope, incoming, send):
+ self.assertEqual(scope['method'], method)
+ for message in original:
+ await send(copied_message(message))
+ messages = []
+ async def send(message):
+ messages.append(copied_message(message))
+ await SimulationResultCompressionMiddleware(application)(
+ http_scope(method, path, accept_encoding), receive, send)
+ return messages, original
+
+ def assert_complete_gzip(self, messages, expected):
+ encoded = b''.join(m.get('body', b'') for m in response_bodies(messages))
+ self.assertEqual(gzip.decompress(encoded), expected)
+ decoder = zlib.decompressobj(zlib.MAX_WBITS + 16)
+ self.assertEqual(decoder.decompress(encoded) + decoder.flush(), expected)
+ self.assertTrue(decoder.eof, 'The gzip trailer must be complete')
+ self.assertEqual(decoder.unused_data, b'', 'No bytes may follow the gzip stream')
+ self.assertEqual(decoder.unconsumed_tail, b'')
+ headers = response_headers(messages)
+ self.assertEqual(headers[b'content-encoding'], b'gzip')
+ vary = [part.strip().lower() for part in headers.get(b'vary', b'').split(b',')]
+ self.assertIn(b'accept-encoding', vary)
+ if b'content-length' in headers:
+ self.assertEqual(int(headers[b'content-length']), len(encoded))
+ bodies = response_bodies(messages)
+ self.assertTrue(bodies)
+ self.assertFalse(bodies[-1].get('more_body', False))
+ self.assertTrue(all(m.get('more_body', False) for m in bodies[:-1]))
+
+ async def test_progress_is_decodable_before_application_sends_result(self):
+ progress = '{"event":"progress","message":"正在计算雪的温度","progress":20}\n'.encode()
+ result = ('{"event":"result","series":{"x":[5e-324,-0,-0.0,1e-300]},'
+ '"message":"温度与压力\\n已完成","padding":"' + 'x' * (CHUNK_BYTES + 200) + '"}\n').encode()
+ decoder = zlib.decompressobj(zlib.MAX_WBITS + 16)
+ received = bytearray()
+ messages = []
+ result_started = False
+ async def send(message):
+ messages.append(copied_message(message))
+ if message['type'] == 'http.response.body':
+ received.extend(decoder.decompress(message.get('body', b'')))
+ if not result_started:
+ self.assertFalse(decoder.eof)
+ async def application(scope, incoming, send):
+ await send({'type': 'http.response.start', 'status': 200, 'headers': [
+ (b'content-type', b'application/x-ndjson'),
+ (b'cache-control', b'no-cache, no-transform'),
+ (b'x-accel-buffering', b'no'), (b'vary', b'Origin'),
+ ]})
+ await send({'type': 'http.response.body', 'body': progress, 'more_body': True})
+ self.assertEqual(bytes(received), progress,
+ 'Progress must be readable before the app computes/sends its result')
+ nonlocal result_started
+ result_started = True
+ # Arbitrary UTF-8/JSON boundaries must not change a byte.
+ utf8_split = result.index('温度'.encode()) + 1
+ for body in (result[:7], result[7:utf8_split], result[utf8_split:]):
+ await send({'type': 'http.response.body', 'body': body, 'more_body': True})
+ await send({'type': 'http.response.body', 'body': b'', 'more_body': False})
+ await SimulationResultCompressionMiddleware(application)(http_scope(), receive, send)
+ self.assertEqual(bytes(received), progress + result)
+ self.assertTrue(decoder.eof)
+ self.assert_complete_gzip(messages, progress + result)
+ headers = response_headers(messages)
+ self.assertNotIn(b'content-length', headers)
+ self.assertEqual(headers[b'cache-control'], b'no-cache, no-transform')
+ self.assertEqual(headers[b'x-accel-buffering'], b'no')
+ self.assertIn(b'origin', headers[b'vary'].lower())
+
+ async def test_large_single_body_is_split_and_preserves_entire_trailer(self):
+ # Deterministic high-entropy ASCII avoids a trivially tiny compressed
+ # result hiding an accidental one-shot write of the whole large body.
+ body = b'{"opaque":"' + hashlib.shake_256(b'asgi-result').hexdigest(600 * 1024).encode() + b'"}'
+ messages, _ = await self.request(body)
+ self.assert_complete_gzip(messages, body)
+ chunks = [m['body'] for m in response_bodies(messages) if m.get('body')]
+ self.assertGreaterEqual(len(chunks), 4)
+ # DEFLATE may enlarge incompressible input slightly, so compressed
+ # messages are allowed framing overhead beyond the 256 KiB input cap.
+ self.assertLessEqual(max(map(len, chunks)), CHUNK_BYTES + 4096)
+ self.assertNotIn(b'content-length', response_headers(messages))
+
+ async def test_gzip_negotiation_and_supported_result_routes(self):
+ body = b'{"result":"' + b'large-result,' * 100 + b'"}'
+ for method, path in [('POST', STREAM_PATH), ('POST', '/api/system-xml/simulate'),
+ ('GET', '/api/system-xml/simulations/task-123')]:
+ for encoding in ('gzip', 'br, gzip;q=0.5', '*', '*;q=0.5', 'GZip;Q=0.7',
+ 'gzip;q=0.5, *;q=0', 'gzip;q=0.8, identity;q=0.5',
+ ['br', 'gzip;q=0.5']):
+ with self.subTest(method=method, path=path, encoding=encoding):
+ messages, _ = await self.request(body, method=method, path=path,
+ accept_encoding=encoding)
+ self.assert_complete_gzip(messages, body)
+
+ async def test_nonaccepting_clients_and_unrelated_routes_are_unchanged(self):
+ body = b'{"result":"' + b'x' * 2000 + b'"}'
+ for encoding in (None, '', 'identity', 'br', 'gzip;q=0', 'gzip;q=0.000',
+ '*;q=0', 'gzip;q=0, *;q=1', 'GZIP;Q=0', 'notgzip',
+ 'gzip;q=bogus', 'gzip;q=1.001', 'gzip;q=0.1234',
+ 'gzip;q=0.8, identity;q=1', ['gzip', 'gzip;q=0']):
+ with self.subTest(encoding=encoding):
+ messages, original = await self.request(body, accept_encoding=encoding)
+ # Identity negotiation still varies by Accept-Encoding for
+ # shared caches; its status, existing headers and bytes stay exact.
+ headers = response_headers(messages)
+ self.assertNotIn(b'content-encoding', headers)
+ self.assertIn(b'accept-encoding', headers[b'vary'].lower())
+ normalized = [copied_message(message) for message in messages]
+ start = next(message for message in normalized if message['type'] == 'http.response.start')
+ start['headers'] = [(key, value) for key, value in start['headers'] if key.lower() != b'vary']
+ self.assertEqual(normalized, original)
+ for method, path in [('GET', STREAM_PATH), ('POST', '/api/system-xml/validate'),
+ ('POST', '/api/system-xml/simulations/task/cancel'),
+ ('GET', '/api/system-xml/simulations/task/cancel'),
+ ('GET', '/api/system-xml/simulations'), ('GET', '/')]:
+ with self.subTest(method=method, path=path):
+ messages, original = await self.request(body, method=method, path=path)
+ self.assertEqual(messages, original)
+
+ async def test_small_preencoded_and_partial_responses_are_not_compressed(self):
+ tiny = b'{"status":"queued"}'
+ messages, original = await self.request(tiny)
+ self.assertEqual(messages, original)
+ body = b'{"data":"' + b'x' * (CHUNK_BYTES + 50) + b'"}'
+ for status, extra in [(200, [(b'content-encoding', b'br')]),
+ (200, [(b'content-encoding', b'gzip')]),
+ (206, [(b'content-range', b'bytes 0-100/2000')])]:
+ with self.subTest(status=status, extra=extra):
+ headers = [(b'content-type', b'application/json'),
+ (b'content-length', str(len(body)).encode()), *extra]
+ messages, original = await self.request(body, status=status, headers=headers)
+ self.assertEqual(messages, original)
+
+ async def test_error_and_cancelled_results_preserve_status_and_numeric_text(self):
+ for code, status in [(422, 'failed'), (200, 'stopped'), (200, 'stalled')]:
+ with self.subTest(code=code, status=status):
+ body = ('{"status":"' + status + '","partial":true,"series":'
+ '{"温度":[-0,5e-324,-0.0,1.7976931348623157e308]},'
+ '"message":"' + '保留已接受的部分结果。' * 80 + '"}').encode()
+ messages, _ = await self.request(body, status=code)
+ self.assertEqual(next(m['status'] for m in messages if m['type'] == 'http.response.start'), code)
+ self.assert_complete_gzip(messages, body)
+
+ async def test_application_failures_and_cancellation_are_not_swallowed(self):
+ progress = b'{"event":"progress","progress":20}\n'
+ for exception in (RuntimeError('upstream failed'), asyncio.CancelledError()):
+ with self.subTest(exception=type(exception).__name__):
+ messages = []
+ async def send(message):
+ messages.append(copied_message(message))
+ async def application(scope, incoming, send):
+ await send({'type': 'http.response.start', 'status': 200,
+ 'headers': [(b'content-type', b'application/x-ndjson')]})
+ await send({'type': 'http.response.body', 'body': progress, 'more_body': True})
+ raise exception
+ with self.assertRaises(type(exception)) as caught:
+ await SimulationResultCompressionMiddleware(application)(http_scope(), receive, send)
+ self.assertIs(caught.exception, exception)
+ decoder = zlib.decompressobj(zlib.MAX_WBITS + 16)
+ partial = b''.join(m.get('body', b'') for m in response_bodies(messages))
+ self.assertEqual(decoder.decompress(partial), progress)
+ self.assertFalse(decoder.eof, 'A failed stream must not masquerade as normally completed')
+
+ async def test_downstream_disconnect_propagates_and_unwinds_application(self):
+ unwound = False
+ failure = BrokenPipeError('client disconnected')
+ async def application(scope, incoming, send):
+ nonlocal unwound
+ try:
+ await send({'type': 'http.response.start', 'status': 200,
+ 'headers': [(b'content-type', b'application/x-ndjson')]})
+ await send({'type': 'http.response.body', 'body': b'{"event":"progress"}\n',
+ 'more_body': True})
+ self.fail('The application must observe the failed downstream send')
+ finally:
+ unwound = True
+ async def send(message):
+ if message['type'] == 'http.response.body':
+ raise failure
+ with self.assertRaises(BrokenPipeError) as caught:
+ await SimulationResultCompressionMiddleware(application)(http_scope(), receive, send)
+ self.assertIs(caught.exception, failure)
+ self.assertTrue(unwound)
+
+
+if __name__ == '__main__':
+ unittest.main()
--- /dev/null
+++ b/tests/manual/result_transfer_profile.py
@@ -0,0 +1,250 @@
+"""Serve a real app with optional throttling AFTER its content compression.
+
+ .venv/bin/python tests/manual/result_transfer_profile.py \
+ --source-root . --frontend-dist frontend/dist --port 8012 \
+ --output-dir test/result-transfer/new --download-bytes-per-second 1600000
+
+Only /api/system-xml/simulate-stream response bodies are split into <=64 KiB
+chunks. Each encoded chunk is delayed before the original ASGI send using
+max(now, previous_deadline) + chunk_bytes/rate; solver idle time earns no credit.
+Rate 0 preserves the same chunking without deliberate delays. No solver hooks,
+request-body copies, response parsing, browser launch, or compilation occurs.
+"""
+from __future__ import annotations
+
+import argparse
+import asyncio
+from hashlib import sha256
+import importlib
+import json
+import logging
+import math
+import os
+from pathlib import Path
+import re
+import shutil
+import sys
+import time
+from uuid import uuid4
+
+ROOT = Path(__file__).resolve().parents[2]
+STREAM_PATH = "/api/system-xml/simulate-stream"
+CHUNK_BYTES = 64 * 1024
+
+
+def write_json(path: Path, value: dict) -> None:
+ path.write_text(json.dumps(value, ensure_ascii=False, indent=2, allow_nan=False) + "\n", encoding="utf-8")
+
+
+def digest(path: Path) -> str:
+ return sha256(path.read_bytes()).hexdigest()
+
+
+class TransferProfile:
+ """An outer ASGI wrapper; outgoing body bytes are already content-encoded."""
+
+ def __init__(self, app, output: Path, rate: int, *, clock=time.perf_counter, sleep=asyncio.sleep):
+ if type(rate) is not int or rate < 0:
+ raise ValueError("download bytes per second must be a nonnegative integer")
+ self.app = app
+ self.output = output
+ self.rate = rate
+ self.clock = clock
+ self.sleep = sleep
+
+ async def __call__(self, scope, receive, send):
+ if scope["type"] != "http" or scope.get("path") != STREAM_PATH:
+ return await self.app(scope, receive, send)
+ request_headers = dict(scope.get("headers", []))
+ sid = request_headers.get(b"x-simulation-id", b"").decode("latin-1")
+ request_id = sid if re.fullmatch(r"[A-Za-z0-9_-]{1,128}", sid) else uuid4().hex
+ directory = self.output / "requests" / request_id
+ try:
+ directory.mkdir(parents=True, exist_ok=False)
+ except FileExistsError:
+ request_id += "-" + uuid4().hex
+ directory = self.output / "requests" / request_id
+ directory.mkdir(parents=True, exist_ok=False)
+ started = self.clock()
+ deadline = started
+ record = {
+ "schemaVersion": 1, "requestId": request_id, "simulationId": sid or None,
+ "path": scope["path"], "method": scope.get("method"), "httpVersion": scope.get("http_version"),
+ "asgiSpecVersion": scope.get("asgi", {}).get("spec_version"),
+ "downloadBytesPerSecond": self.rate, "maxChunkBytes": CHUNK_BYTES,
+ "requestAcceptEncoding": request_headers.get(b"accept-encoding", b"").decode("latin-1"),
+ "status": "incomplete", "responseHeaders": [], "httpStatus": None,
+ "encodedBodyBytes": 0, "attemptedEncodedBodyBytes": 0,
+ "sourceBodyMessages": 0, "chunkCount": 0, "nonEmptyChunkCount": 0,
+ "maxObservedChunkBytes": 0, "throttleSleepSeconds": 0.0, "sendAwaitSeconds": 0.0,
+ "responseHeadersSeconds": None, "firstBodyAvailableSeconds": None,
+ "firstNonemptyBodyAvailableSeconds": None, "firstBodySendStartSeconds": None,
+ "firstBodySendEndSeconds": None, "lastBodySendStartSeconds": None,
+ "lastBodySendEndSeconds": None, "responseBodyCompleteSeconds": None,
+ "disconnectObserved": False, "disconnectBeforeBodyComplete": False,
+ "disconnectObservedSeconds": None, "sendDisconnectException": None,
+ "applicationReturned": False, "exception": None,
+ }
+ body_complete = False
+ active_exception = None
+
+ def elapsed() -> float:
+ return self.clock() - started
+
+ async def observed_receive():
+ message = await receive()
+ if message["type"] == "http.disconnect":
+ record["disconnectObserved"] = True
+ if record["disconnectObservedSeconds"] is None:
+ record["disconnectObservedSeconds"] = elapsed()
+ # Uvicorn may report http.disconnect after a normally finished response.
+ if not body_complete:
+ record["disconnectBeforeBodyComplete"] = True
+ return message
+
+ async def send_chunk(message, chunk: bytes, more_body: bool):
+ nonlocal deadline, body_complete
+ length = len(chunk)
+ if self.rate and length:
+ now = self.clock()
+ deadline = max(now, deadline) + length / self.rate
+ pause_start = self.clock()
+ try:
+ await self.sleep(max(0.0, deadline - pause_start))
+ finally:
+ record["throttleSleepSeconds"] += self.clock() - pause_start
+ start = elapsed()
+ if length and record["firstBodySendStartSeconds"] is None:
+ record["firstBodySendStartSeconds"] = start
+ record["lastBodySendStartSeconds"] = start
+ record["attemptedEncodedBodyBytes"] += length
+ try:
+ await send({**message, "body": chunk, "more_body": more_body})
+ except OSError as exc:
+ record["sendDisconnectException"] = {"type": type(exc).__name__, "message": str(exc)[:1000]}
+ record["disconnectBeforeBodyComplete"] = not body_complete
+ raise
+ finally:
+ record["sendAwaitSeconds"] += elapsed() - start
+ end = elapsed()
+ record["encodedBodyBytes"] += length
+ record["chunkCount"] += 1
+ record["nonEmptyChunkCount"] += int(length > 0)
+ record["maxObservedChunkBytes"] = max(record["maxObservedChunkBytes"], length)
+ if length and record["firstBodySendEndSeconds"] is None:
+ record["firstBodySendEndSeconds"] = end
+ record["lastBodySendEndSeconds"] = end
+ if not more_body:
+ body_complete = True
+ record["responseBodyCompleteSeconds"] = end
+
+ async def observed_send(message):
+ if message["type"] == "http.response.start":
+ record["responseHeadersSeconds"] = elapsed()
+ record["httpStatus"] = message["status"]
+ record["responseHeaders"] = [[key.decode("latin-1"), value.decode("latin-1")]
+ for key, value in message.get("headers", [])]
+ headers = dict(message.get("headers", []))
+ record["contentEncoding"] = headers.get(b"content-encoding", b"").decode("latin-1") or None
+ record["contentType"] = headers.get(b"content-type", b"").decode("latin-1") or None
+ return await send(message)
+ if message["type"] != "http.response.body":
+ return await send(message)
+ body = message.get("body", b"")
+ record["sourceBodyMessages"] += 1
+ if record["firstBodyAvailableSeconds"] is None:
+ record["firstBodyAvailableSeconds"] = elapsed()
+ if body and record["firstNonemptyBodyAvailableSeconds"] is None:
+ record["firstNonemptyBodyAvailableSeconds"] = elapsed()
+ more_body = message.get("more_body", False)
+ if not body:
+ await send_chunk(message, b"", more_body)
+ return
+ for offset in range(0, len(body), CHUNK_BYTES):
+ end = min(offset + CHUNK_BYTES, len(body))
+ await send_chunk(message, body[offset:end], more_body or end < len(body))
+
+ try:
+ await self.app(scope, observed_receive, observed_send)
+ record["applicationReturned"] = True
+ except BaseException as exc:
+ active_exception = exc
+ record["exception"] = {"type": type(exc).__name__, "message": str(exc)[:1000]}
+ raise # Preserve disconnect, cancellation and application failures.
+ finally:
+ record["responseSeconds"] = elapsed()
+ record["responseBodyComplete"] = body_complete
+ if record["disconnectBeforeBodyComplete"] or record["sendDisconnectException"]:
+ record["status"] = "disconnected"
+ elif active_exception is not None:
+ record["status"] = "cancelled" if isinstance(active_exception, asyncio.CancelledError) else "failed"
+ elif body_complete:
+ record["status"] = "completed"
+ try:
+ write_json(directory / "transfer.json", record)
+ except OSError:
+ if active_exception is None:
+ raise
+ logging.exception("Could not save transfer metadata while propagating the original exception")
+
+
+def main() -> None:
+ parser = argparse.ArgumentParser(description=__doc__, formatter_class=argparse.RawDescriptionHelpFormatter)
+ parser.add_argument("--source-root", type=Path, default=ROOT, help="Repository or frozen source snapshot containing app/main.py")
+ parser.add_argument("--frontend-dist", type=Path, help="Default: SOURCE_ROOT/frontend/dist; copied to the fresh output directory")
+ parser.add_argument("--port", type=int, default=8012)
+ parser.add_argument("--output-dir", type=Path, required=True)
+ parser.add_argument("--download-bytes-per-second", type=int, default=0, help="Encoded HTTP body bytes/s; 0 disables deliberate delays")
+ args = parser.parse_args()
+ source, output = args.source_root.resolve(), args.output_dir.resolve()
+ frontend = (args.frontend_dist or source / "frontend/dist").resolve()
+ if not (source / "app/main.py").is_file():
+ parser.error("--source-root must contain app/main.py")
+ if not (frontend / "index.html").is_file():
+ parser.error("--frontend-dist must contain index.html")
+ if output.exists() or output.is_relative_to(frontend):
+ parser.error("Choose a fresh output directory outside --frontend-dist")
+ if not 1 <= args.port <= 65535 or args.download_bytes_per_second < 0:
+ parser.error("Require port 1..65535 and download bytes per second >= 0")
+ if not math.isfinite(float(args.download_bytes_per_second)):
+ parser.error("Download rate must be finite")
+ output.mkdir(parents=True)
+ shutil.copytree(frontend, output / "frontend")
+ # These timing-only inherited switches must not accidentally instrument a run.
+ for name in ("SIMULATIONAPP_PROFILE", "NATIVE_COMPUTE_PROFILE", "NATIVE_STAGE_PROFILE"):
+ os.environ.pop(name, None)
+ sys.path.insert(0, str(source))
+ api = importlib.import_module("app.main")
+ if Path(api.__file__).resolve() != source / "app/main.py":
+ raise RuntimeError("Imported app.main does not belong to the requested source root")
+ api.FRONTEND_DIST_DIR = output / "frontend"
+ import uvicorn
+
+ metadata = {
+ "schemaVersion": 1, "mode": "encoded-body-rate-limit", "sourceRoot": str(source),
+ "loadedApp": str(Path(api.__file__).resolve()), "appMainSha256": digest(source / "app/main.py"),
+ "scriptSha256": digest(Path(__file__)), "frontendDist": str(frontend),
+ "frontendFiles": {str(p.relative_to(output / "frontend")): digest(p)
+ for p in sorted((output / "frontend").rglob("*")) if p.is_file()},
+ "python": sys.version, "platform": sys.platform, "uvicorn": uvicorn.__version__,
+ "host": "127.0.0.1", "port": args.port,
+ "downloadBytesPerSecond": args.download_bytes_per_second, "maxChunkBytes": CHUNK_BYTES,
+ "limitedPath": STREAM_PATH, "solverInstrumented": False,
+ "definitions": {
+ "rate": "Each nonempty encoded chunk is delayed before ASGI send with deadline=max(now,last_deadline)+chunk_bytes/rate; no accumulated credit during computation or socket backpressure. 0 means no deliberate delay, with the same <=64KiB chunking.",
+ "encodedBodyBytes": "Body bytes after application content compression whose original ASGI send returned successfully. Includes progress/result/newline bytes; excludes HTTP framing, response headers, TCP/TLS and port-forward overhead. Not proof of client receipt.",
+ "chunkCount": "Successfully sent ASGI body messages, including empty messages. nonEmptyChunkCount excludes them.",
+ "bodyTimes": "Seconds since this request entered the outer middleware. firstBodySend* refers to first nonempty chunk; lastBodySend* includes an empty terminal body. responseSeconds ends when the application returns/raises, before metadata file I/O.",
+ "completion": "completed means the final more_body=false ASGI send returned and the app returned normally. It does not mean the browser parsed or rendered the result.",
+ "disconnect": "Original receive messages and raised exceptions are preserved. Disconnect observed after a completed body is separately recorded but does not relabel a normal response. Before-completion disconnect, cancellation and failure are distinct statuses.",
+ "scope": "Only the simulation streaming route is throttled. This harness models encoded HTTP body throughput, not latency, packet loss or an actual remote tunnel. Current/frozen applications retain their own compression policy.",
+ },
+ }
+ write_json(output / "environment.json", metadata)
+ print(json.dumps({"url": f"http://127.0.0.1:{args.port}", "output": str(output),
+ "downloadBytesPerSecond": args.download_bytes_per_second, "sourceRoot": str(source)}), flush=True)
+ uvicorn.run(TransferProfile(api.app, output, args.download_bytes_per_second), host="127.0.0.1", port=args.port)
+
+
+if __name__ == "__main__":
+ main()