|
| 1 | +//! web 扩展收件命令面(v0.20.4 / REQ-304 阶段 2 薄壳——本地服务起停/状态)。 |
| 2 | +//! |
| 3 | +//! @ai-context: 服务=std TcpListener 单线程小循环(只绑 127.0.0.1 随机端口、 |
| 4 | +//! 随机 token 首启生成入 data_dir/web_inbox.json、/ping 探测 |
| 5 | +//! (Joplin 范式)、POST /ingest 单向投递、OPTIONS 预检放行、 |
| 6 | +//! CORS 仅允许扩展用途(Authorization 头)+ token 校验 401); |
| 7 | +//! 投递成功=建 kind=web 会话+页面(与 URL 采集同收口),图 base64 |
| 8 | +//! 落盘 notes-images/ 并改写 md 引用。 |
| 9 | +//! @ai-context: 安全边界:token 仅本机展示/存盘(权限收紧为当前用户可读写)、 |
| 10 | +//! 载荷白名单校验(web_inbox::validate_payload)、体量上限; |
| 11 | +//! 服务仅应用运行期有效(进程退出即停)——无驻留后门。 |
| 12 | +
|
| 13 | +use serde::{Deserialize, Serialize}; |
| 14 | +use std::collections::HashMap; |
| 15 | +use std::io::{Read, Write}; |
| 16 | +use std::net::{TcpListener, TcpStream}; |
| 17 | +use std::sync::atomic::{AtomicBool, Ordering}; |
| 18 | +use std::sync::{Arc, Mutex}; |
| 19 | +use std::time::{SystemTime, UNIX_EPOCH}; |
| 20 | +use tauri::State; |
| 21 | + |
| 22 | +use crate::commands::AppState; |
| 23 | +use crate::db_web::WebPage; |
| 24 | +use crate::web_capture::host_of; |
| 25 | +use crate::web_inbox::{is_authorized, parse_headers, validate_payload, IngestPayload}; |
| 26 | +/// 投递体量上限(头+体,防慢速拖垮)。 |
| 27 | +const BODY_MAX: usize = 8 * 1024 * 1024; |
| 28 | + |
| 29 | +/// 收件服务运行时(内存态;token 同时持久化供重启复用)。 |
| 30 | +#[derive(Clone)] |
| 31 | +pub struct WebInboxRuntime { |
| 32 | + pub port: u16, |
| 33 | + pub token: String, |
| 34 | + stop: Arc<AtomicBool>, |
| 35 | +} |
| 36 | + |
| 37 | +impl WebInboxRuntime { |
| 38 | + fn new(port: u16, token: String) -> Self { |
| 39 | + Self { port, token, stop: Arc::new(AtomicBool::new(false)) } |
| 40 | + } |
| 41 | + pub fn request_stop(&self) { |
| 42 | + self.stop.store(true, Ordering::SeqCst); |
| 43 | + } |
| 44 | + pub fn stopping(&self) -> bool { |
| 45 | + self.stop.load(Ordering::SeqCst) |
| 46 | + } |
| 47 | +} |
| 48 | + |
| 49 | +/// 状态视图(设置页/课堂助手展示:端口 + token 复制给扩展)。 |
| 50 | +#[derive(Debug, Clone, Serialize)] |
| 51 | +#[serde(rename_all = "camelCase")] |
| 52 | +pub struct WebInboxView { |
| 53 | + pub running: bool, |
| 54 | + pub port: Option<u16>, |
| 55 | + pub token: Option<String>, |
| 56 | + pub inbox_url: Option<String>, |
| 57 | +} |
| 58 | + |
| 59 | +fn token_file(data_dir: &std::path::Path) -> std::path::PathBuf { |
| 60 | + data_dir.join("web_inbox.json") |
| 61 | +} |
| 62 | + |
| 63 | +fn load_or_make_token(data_dir: &std::path::Path) -> String { |
| 64 | + if let Ok(raw) = std::fs::read_to_string(token_file(data_dir)) { |
| 65 | + if let Ok(v) = serde_json::from_str::<serde_json::Value>(&raw) { |
| 66 | + if let Some(t) = v.get("token").and_then(|x| x.as_str()) { |
| 67 | + if t.len() == 24 { |
| 68 | + return t.to_string(); |
| 69 | + } |
| 70 | + } |
| 71 | + } |
| 72 | + } |
| 73 | + let seed = SystemTime::now() |
| 74 | + .duration_since(UNIX_EPOCH) |
| 75 | + .map(|d| d.as_nanos() as u64) |
| 76 | + .unwrap_or(0) |
| 77 | + ^ std::process::id() as u64; |
| 78 | + let token = crate::web_inbox::random_token(seed); |
| 79 | + if let Ok(raw) = serde_json::to_string_pretty(&serde_json::json!({ "token": token })) { |
| 80 | + let _ = std::fs::write(token_file(data_dir), raw); |
| 81 | + } |
| 82 | + token |
| 83 | +} |
| 84 | + |
| 85 | +/// 启动收件服务(幂等:已运行返回现状)。 |
| 86 | +#[tauri::command] |
| 87 | +pub fn web_inbox_start(state: State<'_, AppState>) -> Result<WebInboxView, String> { |
| 88 | + let mut slot = state |
| 89 | + .web_inbox |
| 90 | + .lock() |
| 91 | + .map_err(|_| "收件服务锁中毒".to_string())?; |
| 92 | + if let Some(rt) = slot.as_ref() { |
| 93 | + return Ok(view_of(rt)); |
| 94 | + } |
| 95 | + let token = load_or_make_token(&state.data_dir); |
| 96 | + let listener = TcpListener::bind("127.0.0.1:0").map_err(|e| format!("绑定回环端口失败: {}", e))?; |
| 97 | + let port = listener.local_addr().map_err(|e| e.to_string())?.port(); |
| 98 | + listener |
| 99 | + .set_nonblocking(true) |
| 100 | + .map_err(|e| format!("设置非阻塞失败: {}", e))?; |
| 101 | + let rt = WebInboxRuntime::new(port, token); |
| 102 | + let runtime = rt.clone(); |
| 103 | + let db = state.db.clone(); |
| 104 | + let data_dir = state.data_dir.clone(); |
| 105 | + let app = state.app.clone(); |
| 106 | + std::thread::Builder::new() |
| 107 | + .name("entropy-web-inbox".into()) |
| 108 | + .spawn(move || loop { |
| 109 | + if runtime.stopping() { |
| 110 | + break; |
| 111 | + } |
| 112 | + match listener.accept() { |
| 113 | + Ok((stream, _)) => handle_connection(stream, &runtime, &db, &data_dir, &app), |
| 114 | + Err(ref e) if e.kind() == std::io::ErrorKind::WouldBlock => { |
| 115 | + std::thread::sleep(std::time::Duration::from_millis(50)); |
| 116 | + } |
| 117 | + Err(e) => { |
| 118 | + eprintln!("[web-inbox] accept 失败: {e}"); |
| 119 | + std::thread::sleep(std::time::Duration::from_millis(200)); |
| 120 | + } |
| 121 | + } |
| 122 | + }) |
| 123 | + .map_err(|e| format!("收件线程启动失败: {e}"))?; |
| 124 | + *slot = Some(rt.clone()); |
| 125 | + Ok(view_of(&rt)) |
| 126 | +} |
| 127 | + |
| 128 | +fn view_of(rt: &WebInboxRuntime) -> WebInboxView { |
| 129 | + WebInboxView { |
| 130 | + running: true, |
| 131 | + port: Some(rt.port), |
| 132 | + token: Some(rt.token.clone()), |
| 133 | + inbox_url: Some(format!("http://127.0.0.1:{}/", rt.port)), |
| 134 | + } |
| 135 | +} |
| 136 | + |
| 137 | +/// 状态查询。 |
| 138 | +#[tauri::command] |
| 139 | +pub fn web_inbox_status(state: State<'_, AppState>) -> Result<WebInboxView, String> { |
| 140 | + let slot = state.web_inbox.lock().map_err(|_| "收件服务锁中毒".to_string())?; |
| 141 | + match slot.as_ref() { |
| 142 | + Some(rt) => Ok(view_of(rt)), |
| 143 | + None => Ok(WebInboxView { running: false, port: None, token: None, inbox_url: None }), |
| 144 | + } |
| 145 | +} |
| 146 | + |
| 147 | +/// 停止服务(token 保留——重启同 token,扩展零重配)。 |
| 148 | +#[tauri::command] |
| 149 | +pub fn web_inbox_stop(state: State<'_, AppState>) -> Result<(), String> { |
| 150 | + let mut slot = state.web_inbox.lock().map_err(|_| "收件服务锁中毒".to_string())?; |
| 151 | + if let Some(rt) = slot.take() { |
| 152 | + rt.request_stop(); |
| 153 | + } |
| 154 | + Ok(()) |
| 155 | +} |
| 156 | + |
| 157 | +fn write_simple(stream: &mut TcpStream, status: &str, body: &str, cors: bool) { |
| 158 | + let mut headers = format!( |
| 159 | + "HTTP/1.1 {}\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n", |
| 160 | + status, |
| 161 | + body.len() |
| 162 | + ); |
| 163 | + if cors { |
| 164 | + headers.push_str("Access-Control-Allow-Origin: *\r\nAccess-Control-Allow-Headers: authorization, content-type\r\nAccess-Control-Allow-Methods: POST, OPTIONS\r\n"); |
| 165 | + } |
| 166 | + let _ = stream.write_all(headers.as_bytes()); |
| 167 | + let _ = stream.write_all(body.as_bytes()); |
| 168 | + let _ = stream.flush(); |
| 169 | +} |
| 170 | + |
| 171 | +fn read_request(stream: &mut TcpStream) -> Option<(String, String, HashMap<String, String>, Vec<u8>)> { |
| 172 | + let mut buf = Vec::new(); |
| 173 | + let mut chunk = [0u8; 4096]; |
| 174 | + let mut head_end: Option<usize> = None; |
| 175 | + loop { |
| 176 | + match stream.read(&mut chunk) { |
| 177 | + Ok(0) => break, |
| 178 | + Ok(n) => { |
| 179 | + buf.extend_from_slice(&chunk[..n]); |
| 180 | + if head_end.is_none() { |
| 181 | + if let Some(pos) = find_head_end(&buf) { |
| 182 | + head_end = Some(pos); |
| 183 | + let content_len: usize = String::from_utf8_lossy(&buf[..pos]) |
| 184 | + .to_ascii_lowercase() |
| 185 | + .lines() |
| 186 | + .find_map(|l| { |
| 187 | + l.trim() |
| 188 | + .split_once(':') |
| 189 | + .filter(|(k, _)| k.trim() == "content-length") |
| 190 | + .and_then(|(_, v)| v.trim().parse().ok()) |
| 191 | + }) |
| 192 | + .unwrap_or(0); |
| 193 | + if pos + 4 + content_len <= buf.len() { |
| 194 | + let body = buf[pos + 4..pos + 4 + content_len].to_vec(); |
| 195 | + let head = String::from_utf8_lossy(&buf[..pos]).into_owned(); |
| 196 | + let (method, path, headers) = parse_headers(&head)?; |
| 197 | + return Some((method, path, headers, body)); |
| 198 | + } |
| 199 | + if buf.len() > BODY_MAX { |
| 200 | + return None; |
| 201 | + } |
| 202 | + } |
| 203 | + } |
| 204 | + if buf.len() > BODY_MAX && head_end.is_none() { |
| 205 | + return None; |
| 206 | + } |
| 207 | + } |
| 208 | + Err(ref e) if e.kind() == std::io::ErrorKind::WouldBlock => { |
| 209 | + std::thread::sleep(std::time::Duration::from_millis(5)); |
| 210 | + } |
| 211 | + Err(_) => break, |
| 212 | + } |
| 213 | + } |
| 214 | + None |
| 215 | +} |
| 216 | + |
| 217 | +fn find_head_end(buf: &[u8]) -> Option<usize> { |
| 218 | + buf.windows(4).position(|w| w == b"\r\n\r\n") |
| 219 | +} |
| 220 | + |
| 221 | +fn handle_connection( |
| 222 | + mut stream: TcpStream, |
| 223 | + runtime: &WebInboxRuntime, |
| 224 | + db: &crate::db::Db, |
| 225 | + data_dir: &std::path::Path, |
| 226 | + app: &tauri::AppHandle, |
| 227 | +) { |
| 228 | + let Some((method, path, headers, body)) = read_request(&mut stream) else { |
| 229 | + write_simple(&mut stream, "400 Bad Request", r#"{"error":"bad request"}"#, true); |
| 230 | + return; |
| 231 | + }; |
| 232 | + if method == "OPTIONS" { |
| 233 | + write_simple(&mut stream, "204 No Content", "", true); |
| 234 | + return; |
| 235 | + } |
| 236 | + if !is_authorized(&headers, &runtime.token) { |
| 237 | + write_simple(&mut stream, "401 Unauthorized", r#"{"error":"unauthorized"}"#, true); |
| 238 | + return; |
| 239 | + } |
| 240 | + match (method.as_str(), path.as_str()) { |
| 241 | + ("GET", "/ping") => write_simple(&mut stream, "200 OK", r#"{"ok":true}"#, true), |
| 242 | + ("POST", "/ingest") => { |
| 243 | + let payload: Result<IngestPayload, _> = serde_json::from_slice(&body); |
| 244 | + match payload { |
| 245 | + Ok(p) => match validate_payload(&p) { |
| 246 | + Ok(()) => match ingest_from_extension(db, data_dir, Some(app), &p) { |
| 247 | + Ok(session_id) => write_simple( |
| 248 | + &mut stream, |
| 249 | + "200 OK", |
| 250 | + &serde_json::json!({ "ok": true, "sessionId": session_id }).to_string(), |
| 251 | + true, |
| 252 | + ), |
| 253 | + Err(e) => write_simple( |
| 254 | + &mut stream, |
| 255 | + "500 Internal Server Error", |
| 256 | + &serde_json::json!({ "error": e }).to_string(), |
| 257 | + true, |
| 258 | + ), |
| 259 | + }, |
| 260 | + Err(e) => write_simple( |
| 261 | + &mut stream, |
| 262 | + "422 Unprocessable Entity", |
| 263 | + &serde_json::json!({ "error": e }).to_string(), |
| 264 | + true, |
| 265 | + ), |
| 266 | + }, |
| 267 | + Err(_) => write_simple(&mut stream, "400 Bad Request", r#"{"error":"invalid json"}"#, true), |
| 268 | + } |
| 269 | + } |
| 270 | + _ => write_simple(&mut stream, "404 Not Found", r#"{"error":"not found"}"#, true), |
| 271 | + } |
| 272 | +} |
| 273 | + |
| 274 | +/// 扩展投递 → kind=web 会话 + 页面(与 URL 采集同收口);图 base64 落盘改写。 |
| 275 | +fn ingest_from_extension( |
| 276 | + db: &crate::db::Db, |
| 277 | + data_dir: &std::path::Path, |
| 278 | + app: Option<&tauri::AppHandle>, |
| 279 | + p: &IngestPayload, |
| 280 | +) -> Result<i64, String> { |
| 281 | + let now = crate::db::unix_seconds(); |
| 282 | + let session = db |
| 283 | + .create_session(&crate::types::NewSession { |
| 284 | + title: p |
| 285 | + .title |
| 286 | + .as_deref() |
| 287 | + .map(|t| t.trim()) |
| 288 | + .filter(|t| !t.is_empty()) |
| 289 | + .map(|t| t.chars().take(100).collect()) |
| 290 | + .unwrap_or_else(|| { |
| 291 | + p.url.as_deref().and_then(host_of).unwrap_or_else(|| "网页".to_string()) |
| 292 | + }), |
| 293 | + source_window: p.url.clone(), |
| 294 | + profile: None, |
| 295 | + kind: Some("web".to_string()), |
| 296 | + }) |
| 297 | + .map_err(|e| e.to_string())?; |
| 298 | + // 图落盘:notes-images/ 通用目录(编辑器相对路径解析基座) |
| 299 | + let mut markdown = p.markdown.clone(); |
| 300 | + let notes_images = data_dir.join("notes-images"); |
| 301 | + let _ = std::fs::create_dir_all(¬es_images); |
| 302 | + for img in &p.images { |
| 303 | + if let Some(bytes) = crate::web_inbox::data_uri_bytes(&img.data_base64) { |
| 304 | + let mime = img.data_base64.split_once(';').map(|(m, _)| m).unwrap_or("data:image/png"); |
| 305 | + let ext = mime.rsplit('/').next().unwrap_or("png"); |
| 306 | + let filename = format!("web-{}-{}.{}", session.id, crate::web_inbox::short_hash(&bytes), ext); |
| 307 | + if std::fs::write(notes_images.join(&filename), bytes).is_ok() { |
| 308 | + // 替换 md 中 `` 引用为相对路径(编辑器同解析基座) |
| 309 | + markdown = markdown.replace( |
| 310 | + &format!("]({})", img.data_base64), |
| 311 | + &format!("](notes-images/{})", filename), |
| 312 | + ); |
| 313 | + } |
| 314 | + } |
| 315 | + } |
| 316 | + db.insert_web_page(&WebPage { |
| 317 | + session_id: session.id, |
| 318 | + url: p.url.clone().unwrap_or_else(|| "".to_string()), |
| 319 | + site: p.site.clone(), |
| 320 | + author: p.author.clone(), |
| 321 | + published: None, |
| 322 | + markdown, |
| 323 | + raw_html: None, |
| 324 | + extracted_ok: true, |
| 325 | + fetched_at: now, |
| 326 | + }) |
| 327 | + .map_err(|e| e.to_string())?; |
| 328 | + if let Some(app) = app { |
| 329 | + crate::notify::emit_changed(app, crate::notify::DataDomain::Sessions); |
| 330 | + } |
| 331 | + Ok(session.id) |
| 332 | +} |
| 333 | + |
| 334 | +#[cfg(test)] |
| 335 | +#[path = "commands_web_inbox_tests.rs"] |
| 336 | +mod tests; |
0 commit comments