From ef1fe9e949efd30f4a018abbfb2f7d306dc9668f Mon Sep 17 00:00:00 2001 From: KiriAky 107 Date: Wed, 9 Sep 2026 10:37:49 +0800 Subject: [PATCH] =?UTF-8?q?fix(sidecar):=20=E5=85=B3=E9=97=AD=E6=8F=A1?= =?UTF-8?q?=E6=89=8B=E5=89=8D=E7=9A=84=E4=BB=A3=E7=90=86=E7=AA=97=E5=8F=A3?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- docs/development/OpenNexus生产验收Runner.md | 2 +- frontend/src-tauri/src/core.rs | 76 ++++++++- frontend/src-tauri/src/main.rs | 177 ++++++++++++++++---- scripts/acceptance_cases/a03_transport.py | 118 +++++++++++++ scripts/phase3_acceptance.py | 5 + 5 files changed, 342 insertions(+), 36 deletions(-) create mode 100644 scripts/acceptance_cases/a03_transport.py diff --git a/docs/development/OpenNexus生产验收Runner.md b/docs/development/OpenNexus生产验收Runner.md index b4075c1..78c6dcd 100644 --- a/docs/development/OpenNexus生产验收Runner.md +++ b/docs/development/OpenNexus生产验收Runner.md @@ -24,4 +24,4 @@ python scripts/phase3-production-acceptance.py ` 报告目录包含 `summary.json`、`case-manifest.json`、`junit.xml`、`cases/.json` 和脱敏的 `logs/.log`。摘要记录 commit、各锁文件 SHA-256、配置摘要与已提供安装产物摘要。日志将仓库、数据根、报告根、用户主目录和配置声明的秘密值替换为占位符,并限制为 10 MiB。报告目录必须为空,避免单例复跑覆盖原始证据。 -当前 runner 与失败闭合行为已实现,A-02 Sidecar、B-01/B-02 凭据、D-01 扩展包以及 S-01/S-02 Sync 客户端 driver 已登记;其余 24 个生产验收 ID 尚未登记,运行时会生成 `NOT_IMPLEMENTED` 证据并退出 1。这用于阻止误报,不是这些用例的验收通过。 +当前 runner 与失败闭合行为已实现,A-02/A-03 Sidecar、B-01/B-02 凭据、D-01 扩展包以及 S-01/S-02 Sync 客户端 driver 已登记;其余 23 个生产验收 ID 尚未登记,运行时会生成 `NOT_IMPLEMENTED` 证据并退出 1。这用于阻止误报,不是这些用例的验收通过。 diff --git a/frontend/src-tauri/src/core.rs b/frontend/src-tauri/src/core.rs index 9a1805d..e3dd3be 100644 --- a/frontend/src-tauri/src/core.rs +++ b/frontend/src-tauri/src/core.rs @@ -362,9 +362,11 @@ impl CoreSupervisor { .write_all(&payload) .map_err(|_| "CORE_PIPE_FAILED")?; let (tx, rx) = mpsc::channel(); + let (activation_tx, activation_rx) = mpsc::channel(); let lifetime = session.lifetime.clone(); std::thread::spawn(move || { let mut reader = BufReader::new(stdout); + let mut activated = false; loop { let mut line = Zeroizing::new(Vec::new()); match reader @@ -382,9 +384,24 @@ impl CoreSupervisor { break; }; if message.get("rpc").is_none() { - let _ = tx.send(Ok::, std::io::Error>(line.to_vec())); + if activated + || tx + .send(Ok::, std::io::Error>(line.to_vec())) + .is_err() + { + break; + } + match activation_rx.recv() { + Ok(true) => activated = true, + _ => break, + } continue; } + // A child has no broker authority until its ready frame has + // passed the protocol, identity, generation, and HMAC checks. + if !activated { + break; + } let request_id = message .get("request_id") .and_then(|v| v.as_str()) @@ -421,15 +438,26 @@ impl CoreSupervisor { .map_err(|_| "CORE_READY_TIMEOUT")? .map_err(|_| "CORE_HANDSHAKE_INVALID")?; if line.len() > 16384 { + let _ = activation_tx.send(false); return Err("CORE_HANDSHAKE_INVALID".into()); } - session.port = verify_ready( + let port = verify_ready( &line, &session.secret, &challenge, &session.generation, session.child.id(), - )?; + ); + match port { + Ok(port) => { + session.port = port; + activation_tx.send(true).map_err(|_| "CORE_PIPE_FAILED")?; + } + Err(error) => { + let _ = activation_tx.send(false); + return Err(error); + } + } Ok(session) } } @@ -476,4 +504,46 @@ mod tests { assert!(verify_ready(&payload, &secret, &challenge, &generation, 124).is_err()); assert!(verify_ready(&payload, &secret, &"04".repeat(32), &generation, 123).is_err()); } + + #[test] + fn protocol_incompatibility_rejects_pre_ready_broker_requests() { + use std::sync::atomic::{AtomicUsize, Ordering}; + + let backend = Path::new(env!("CARGO_MANIFEST_DIR")) + .join("../../backend") + .canonicalize() + .unwrap(); + let python = backend.join(if cfg!(windows) { + ".venv/Scripts/python.exe" + } else { + ".venv/bin/python" + }); + assert!(python.is_file(), "backend virtual environment is required"); + let temp = tempfile::tempdir().unwrap(); + let fixture = temp.path().join("incompatible.py"); + std::fs::write( + &fixture, + r#"import json, os, sys, time +config = json.loads(sys.stdin.buffer.readline()) +print(json.dumps({"protocol": 2, "launcher_pid": config["launcher_pid"], "pid": os.getpid(), "port": 1, "generation": config["generation"], "proof": "00" * 32}), flush=True) +print(json.dumps({"rpc": "workspace.write", "request_id": "must-not-run", "params": {}}), flush=True) +time.sleep(5) +"#, + ) + .unwrap(); + let broker_calls = Arc::new(AtomicUsize::new(0)); + let observed = Arc::clone(&broker_calls); + let mut core = CoreSupervisor::new( + python, + vec![fixture.to_string_lossy().into_owned()], + backend, + temp.path().join("data"), + ) + .with_broker(Arc::new(move |_| { + observed.fetch_add(1, Ordering::SeqCst); + Err("BUSINESS_CALL_MUST_NOT_RUN".into()) + })); + assert_eq!(core.start().unwrap_err(), "PROTOCOL_INCOMPATIBLE"); + assert_eq!(broker_calls.load(Ordering::SeqCst), 0); + } } diff --git a/frontend/src-tauri/src/main.rs b/frontend/src-tauri/src/main.rs index 3d8f82f..88b23fa 100644 --- a/frontend/src-tauri/src/main.rs +++ b/frontend/src-tauri/src/main.rs @@ -108,6 +108,49 @@ struct CoreResponse { // 后端常规导出上限为 20 MiB;为 PDF 和未来的二进制接口保留余量,同时限制 IPC 内存占用。 const MAX_CORE_RESPONSE_BYTES: usize = 64 * 1024 * 1024; +fn validate_core_json_body(body: Option<&serde_json::Value>) -> Result<(), String> { + if body.is_some_and(|value| value.to_string().len() > MAX_CORE_RESPONSE_BYTES) { + return Err("CORE_REQUEST_TOO_LARGE".into()); + } + Ok(()) +} + +fn decode_core_binary_body(encoded: String) -> Result, String> { + if encoded.len() > MAX_CORE_RESPONSE_BYTES * 4 / 3 + 4 { + return Err("CORE_REQUEST_TOO_LARGE".into()); + } + let bytes = BASE64_STANDARD + .decode(encoded) + .map_err(|_| "CORE_BODY_INVALID")?; + if bytes.len() > MAX_CORE_RESPONSE_BYTES { + return Err("CORE_REQUEST_TOO_LARGE".into()); + } + Ok(bytes) +} + +fn checked_core_response_size(current: usize, additional: usize) -> Result { + let size = current.saturating_add(additional); + if size > MAX_CORE_RESPONSE_BYTES { + return Err("CORE_RESPONSE_TOO_LARGE".into()); + } + Ok(size) +} + +async fn read_core_response(mut response: reqwest::Response) -> Result, String> { + if response + .content_length() + .is_some_and(|length| length > MAX_CORE_RESPONSE_BYTES as u64) + { + return Err("CORE_RESPONSE_TOO_LARGE".into()); + } + let mut bytes = Vec::new(); + while let Some(chunk) = response.chunk().await.map_err(|_| "CORE_RESPONSE_ERROR")? { + checked_core_response_size(bytes.len(), chunk.len())?; + bytes.extend_from_slice(&chunk); + } + Ok(bytes) +} + fn is_json_content_type(content_type: &str) -> bool { let media_type = content_type .split(';') @@ -125,7 +168,10 @@ fn core_url(path: &str) -> Result { #[cfg(test)] mod core_proxy_tests { - use super::{core_url, is_json_content_type}; + use super::{ + checked_core_response_size, core_url, decode_core_binary_body, is_json_content_type, + read_core_response, validate_core_json_body, MAX_CORE_RESPONSE_BYTES, + }; #[test] fn request_dto_accepts_camel_case_and_rejects_unowned_headers() { @@ -162,6 +208,99 @@ mod core_proxy_tests { "application/vnd.openxmlformats-officedocument.wordprocessingml.document" )); } + + #[tokio::test] + async fn a03_frozen_transfer_limits_and_failure_semantics() { + use std::io::{Read, Write}; + use std::net::TcpListener; + + let exact_json = serde_json::Value::String("x".repeat(MAX_CORE_RESPONSE_BYTES - 2)); + validate_core_json_body(Some(&exact_json)).unwrap(); + drop(exact_json); + let over_json = serde_json::Value::String("x".repeat(MAX_CORE_RESPONSE_BYTES - 1)); + assert_eq!( + validate_core_json_body(Some(&over_json)).unwrap_err(), + "CORE_REQUEST_TOO_LARGE" + ); + drop(over_json); + + let encoded_length = MAX_CORE_RESPONSE_BYTES.div_ceil(3) * 4; + let mut exact_binary = "A".repeat(encoded_length - 2); + exact_binary.push_str("=="); + assert_eq!( + decode_core_binary_body(exact_binary).unwrap().len(), + MAX_CORE_RESPONSE_BYTES + ); + let mut over_binary = "A".repeat(encoded_length - 1); + over_binary.push('='); + assert_eq!( + decode_core_binary_body(over_binary).unwrap_err(), + "CORE_REQUEST_TOO_LARGE" + ); + assert_eq!( + checked_core_response_size(MAX_CORE_RESPONSE_BYTES - 1, 1).unwrap(), + MAX_CORE_RESPONSE_BYTES + ); + assert_eq!( + checked_core_response_size(MAX_CORE_RESPONSE_BYTES, 1).unwrap_err(), + "CORE_RESPONSE_TOO_LARGE" + ); + + fn server(declared: usize, sent: usize) -> (String, std::thread::JoinHandle<()>) { + let listener = TcpListener::bind("127.0.0.1:0").unwrap(); + let url = format!("http://{}/fixture", listener.local_addr().unwrap()); + let worker = std::thread::spawn(move || { + let (mut socket, _) = listener.accept().unwrap(); + let mut request = Vec::new(); + while !request.ends_with(b"\r\n\r\n") { + let mut byte = [0]; + socket.read_exact(&mut byte).unwrap(); + request.push(byte[0]); + } + write!( + socket, + "HTTP/1.1 200 OK\r\nContent-Length: {declared}\r\nContent-Type: application/octet-stream\r\nConnection: close\r\n\r\n" + ) + .unwrap(); + let block = [37u8; 64 * 1024]; + let mut remaining = sent; + while remaining > 0 { + let count = remaining.min(block.len()); + if socket.write_all(&block[..count]).is_err() { + break; + } + remaining -= count; + } + }); + (url, worker) + } + let client = reqwest::Client::builder().no_proxy().build().unwrap(); + let (url, worker) = server(MAX_CORE_RESPONSE_BYTES, MAX_CORE_RESPONSE_BYTES); + let bytes = read_core_response(client.get(url).send().await.unwrap()) + .await + .unwrap(); + assert_eq!(bytes.len(), MAX_CORE_RESPONSE_BYTES); + assert!(bytes.iter().all(|byte| *byte == 37)); + worker.join().unwrap(); + + let (url, worker) = server(MAX_CORE_RESPONSE_BYTES + 1, 0); + assert_eq!( + read_core_response(client.get(url).send().await.unwrap()) + .await + .unwrap_err(), + "CORE_RESPONSE_TOO_LARGE" + ); + worker.join().unwrap(); + + let (url, worker) = server(8, 1); + assert_eq!( + read_core_response(client.get(url).send().await.unwrap()) + .await + .unwrap_err(), + "CORE_RESPONSE_ERROR" + ); + worker.join().unwrap(); + } } /// Authenticated process-local transport; session headers are owned by Rust. @@ -211,12 +350,7 @@ async fn core_request(request: CoreRequest, host: State<'_, Host>) -> Result MAX_CORE_RESPONSE_BYTES) - { - return Err("CORE_REQUEST_TOO_LARGE".into()); - } + validate_core_json_body(body.as_ref())?; let core = host.core.clone(); let core_path = path.clone(); let session = tauri::async_runtime::spawn_blocking(move || { @@ -261,21 +395,13 @@ async fn core_request(request: CoreRequest, host: State<'_, Host>) -> Result MAX_CORE_RESPONSE_BYTES * 4 / 3 + 4 { - return Err("CORE_REQUEST_TOO_LARGE".into()); - } - let bytes = BASE64_STANDARD - .decode(encoded) - .map_err(|_| "CORE_BODY_INVALID")?; - if bytes.len() > MAX_CORE_RESPONSE_BYTES { - return Err("CORE_REQUEST_TOO_LARGE".into()); - } let content_type = content_type .as_deref() .unwrap_or("application/octet-stream"); if !matches!(content_type, "application/octet-stream" | "application/zip") { return Err("CORE_CONTENT_TYPE_DENIED".into()); } + let bytes = decode_core_binary_body(encoded)?; request = request .header(reqwest::header::CONTENT_TYPE, content_type) .body(bytes); @@ -287,7 +413,7 @@ async fn core_request(request: CoreRequest, host: State<'_, Host>) -> Result) -> Result MAX_CORE_RESPONSE_BYTES as u64) - { - return Err("CORE_RESPONSE_TOO_LARGE".into()); - } - let mut bytes = Vec::new(); - while let Some(chunk) = response.chunk().await.map_err(|_| "CORE_RESPONSE_ERROR")? { - if bytes.len().saturating_add(chunk.len()) > MAX_CORE_RESPONSE_BYTES { - return Err("CORE_RESPONSE_TOO_LARGE".into()); - } - bytes.extend_from_slice(&chunk); - } + let bytes = read_core_response(response).await?; let (body, body_base64) = if is_json_content_type(&content_type) { ( String::from_utf8(bytes.to_vec()).map_err(|_| "CORE_RESPONSE_ERROR")?, @@ -396,8 +510,7 @@ fn core_stream( channel.send(serde_json::json!({"kind":"headers","status":response.status().as_u16()})).map_err(|_| "CORE_STREAM_CLOSED")?; let mut size = 0usize; while let Some(bytes) = response.chunk().await.map_err(|_| "CORE_RESPONSE_ERROR")? { - size = size.saturating_add(bytes.len()); - if size > MAX_CORE_RESPONSE_BYTES { return Err("CORE_RESPONSE_TOO_LARGE".into()); } + size = checked_core_response_size(size, bytes.len())?; for chunk in bytes.chunks(16384) { channel.send(serde_json::json!({"kind":"chunk","data":BASE64_STANDARD.encode(chunk)})).map_err(|_| "CORE_STREAM_CLOSED")?; } diff --git a/scripts/acceptance_cases/a03_transport.py b/scripts/acceptance_cases/a03_transport.py new file mode 100644 index 0000000..d4aa6e3 --- /dev/null +++ b/scripts/acceptance_cases/a03_transport.py @@ -0,0 +1,118 @@ +"""A-03 protocol, transport-limit, cancellation, and commit acceptance driver.""" + +from __future__ import annotations + +import argparse +import hashlib +import json +import os +import shutil +import subprocess +from pathlib import Path + +ROOT = Path(__file__).resolve().parents[2] +MANIFEST = ROOT / "frontend" / "src-tauri" / "Cargo.toml" +FRONTEND = ROOT / "frontend" +RUST_TESTS = ( + (("--lib",), "core::tests::protocol_incompatibility_rejects_pre_ready_broker_requests"), + (("--bin", "notesagent-desktop"), "core_proxy_tests::a03_frozen_transfer_limits_and_failure_semantics"), + (("--lib",), "request_lifecycle::tests::cancellation_before_dispatch_and_replay_never_run_work"), + (("--lib",), "request_lifecycle::tests::cancelling_a_real_response_closes_its_socket"), + (("--lib",), "workspace::tests::operation_receipt_survives_reopen_and_replay_after_later_edit"), + (("--test", "core_workspace"), "real_core_notes_roundtrip_only_through_bound_host_and_confirm_commits"), +) +FRONTEND_TESTS = ( + "src/services/apiClient.cancellation.spec.ts", + "src/services/platform/coreStream.spec.ts", +) + + +def sha256(path: Path) -> str: + digest = hashlib.sha256() + with path.open("rb") as source: + for chunk in iter(lambda: source.read(1024 * 1024), b""): + digest.update(chunk) + return digest.hexdigest() + + +def run(command: list[str], cwd: Path, expected: str) -> bool: + completed = subprocess.run(command, cwd=cwd, capture_output=True, text=True, check=False) + print(completed.stdout, end="") + print(completed.stderr, end="") + return completed.returncode == 0 and expected in completed.stdout + + +def main() -> int: + parser = argparse.ArgumentParser() + parser.add_argument("--config", required=True) + parser.add_argument("--output", required=True) + args = parser.parse_args() + case_id = os.environ.get("OPENNEXUS_ACCEPTANCE_CASE_ID", "") + output = Path(args.output) + output.parent.mkdir(parents=True, exist_ok=True) + assertions = [ + {"name": "an incompatible protocol returns PROTOCOL_INCOMPATIBLE before any broker business call"}, + {"name": "JSON and binary requests accept exactly 64 MiB and reject the next byte"}, + {"name": "responses accept exactly 64 MiB while declared oversize and truncated bodies fail closed"}, + {"name": "pre-dispatch cancellation runs no work and in-flight cancellation closes the real socket"}, + {"name": "WebView stream abort, reader failure, and native cancellation retain frozen error semantics"}, + {"name": "a committed mutation remains queryable and exactly replayable by operation_id after reopen"}, + ] + cargo = shutil.which(os.environ.get("CARGO", "cargo")) + pnpm = shutil.which("pnpm.cmd" if os.name == "nt" else "pnpm") + passed = case_id == "A-03" and cargo is not None and pnpm is not None + if passed: + for selector, test in RUST_TESTS: + command = [ + cargo, + "test", + "--manifest-path", + str(MANIFEST), + "--locked", + "--features", + "desktop", + *selector, + test, + "--", + "--exact", + "--nocapture", + ] + passed = run(command, ROOT, "1 passed; 0 failed") and passed + passed = run( + [pnpm, "exec", "vitest", "run", *FRONTEND_TESTS, "--reporter=dot"], + FRONTEND, + "9 passed", + ) and passed + status = "PASSED" if passed else "FAILED" + evidence = "six exact Rust transport/Workspace oracles and two exact frontend transport suites" + for assertion in assertions: + assertion.update({"status": status, "evidence": evidence}) + files = [] + for relative in ( + "frontend/src-tauri/src/core.rs", + "frontend/src-tauri/src/main.rs", + "frontend/src-tauri/src/request_lifecycle.rs", + "frontend/src-tauri/src/workspace.rs", + "frontend/src/services/platform/coreRequest.ts", + "frontend/src/services/platform/coreStream.ts", + ): + files.append({"path": relative, "sha256": sha256(ROOT / relative)}) + payload = { + "schema": 1, + "case_id": case_id, + "status": status, + "reason": "" if passed else "An A-03 protocol, limit, cancellation, or commit oracle failed.", + "assertions": assertions, + "metrics": {"peak_rss_bytes": None, "max_process_count": None, "denied_access_count": None}, + "files": files, + "revisions": [ + {"scope": "frozen transport bytes", "accepted": 67108864, "rejected": 67108865}, + {"scope": "committed operation replay", "replays": 100}, + ], + } + output.write_text(json.dumps(payload, ensure_ascii=False, indent=2) + "\n", encoding="utf-8") + return 0 if passed else 1 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/scripts/phase3_acceptance.py b/scripts/phase3_acceptance.py index 184cfdf..e30497e 100644 --- a/scripts/phase3_acceptance.py +++ b/scripts/phase3_acceptance.py @@ -37,6 +37,11 @@ CASE_DRIVERS: dict[str, dict[str, Any]] = { "timeout_seconds": 900, "required_metrics": (), }, + "A-03": { + "driver": "scripts/acceptance_cases/a03_transport.py", + "timeout_seconds": 900, + "required_metrics": (), + }, "B-01": { "driver": "scripts/acceptance_cases/b01_credentials.py", "timeout_seconds": 900,