use std::collections::{HashMap, HashSet, VecDeque}; use std::net::IpAddr; use std::net::SocketAddr; use std::path::PathBuf; use std::sync::atomic::{AtomicBool, AtomicU64, Ordering}; use std::sync::{Arc, Mutex}; use std::time::{SystemTime, UNIX_EPOCH}; use serde::Serialize; use sha2::{Digest, Sha256}; use uuid::Uuid; use crate::base_system::context::Config; use crate::download::downloader::{BookNameOption, ProgressSnapshot}; #[derive(Clone, Debug)] pub(crate) struct ConfigView { pub(crate) old_cli: bool, pub(crate) use_official_api: bool, pub(crate) save_path: String, pub(crate) api_endpoints_len: usize, } #[derive(Clone)] pub(crate) struct AppState { pub(crate) bind_addrs: Arc>, pub(crate) config_view: Arc, pub(crate) config: Arc>, // allow runtime updates via Web UI pub(crate) config_path: Arc, pub(crate) library_root: Arc, pub(crate) jobs: Arc, pub(crate) self_update: Arc, pub(crate) library_scan: Arc, pub(crate) update_scan: Arc, pub(crate) auth: Option, /// 限制同时访问上游 API(search / preview)的并发数,防止 WebUI 被用作多用户 API 代理。 /// 仅在启用 official-api feature 时有意义,其他 feature 下置 None。 #[cfg(feature = "official-api")] pub(crate) api_semaphore: Arc, } #[derive(Debug, Clone, Copy, Serialize)] #[serde(rename_all = "snake_case")] pub(crate) enum SelfUpdateState { Idle, Running, Done, Failed, } #[derive(Debug, Clone, Serialize)] pub(crate) struct SelfUpdateInfo { pub(crate) state: SelfUpdateState, pub(crate) stage: String, pub(crate) percent: u8, pub(crate) message: String, pub(crate) updated_ms: u64, } #[derive(Debug)] pub(crate) struct SelfUpdateStore { running: AtomicBool, inner: Mutex, } impl Default for SelfUpdateStore { fn default() -> Self { Self { running: AtomicBool::new(false), inner: Mutex::new(SelfUpdateInfo { state: SelfUpdateState::Idle, stage: "idle".to_string(), percent: 0, message: "尚未开始".to_string(), updated_ms: now_ms(), }), } } } impl SelfUpdateStore { pub(crate) fn try_start(&self) -> bool { if self .running .compare_exchange(false, true, Ordering::SeqCst, Ordering::SeqCst) .is_err() { return false; } self.set(SelfUpdateState::Running, "prepare", 2, "准备开始自更新…"); true } pub(crate) fn snapshot(&self) -> SelfUpdateInfo { self.inner.lock().unwrap_or_else(|e| e.into_inner()).clone() } pub(crate) fn set( &self, state: SelfUpdateState, stage: impl Into, percent: u8, message: impl Into, ) { let mut g = self.inner.lock().unwrap_or_else(|e| e.into_inner()); g.state = state; g.stage = stage.into(); g.percent = percent.min(100); g.message = message.into(); g.updated_ms = now_ms(); } pub(crate) fn tick_running(&self) { let mut g = self.inner.lock().unwrap_or_else(|e| e.into_inner()); if !matches!(g.state, SelfUpdateState::Running) { return; } if g.percent < 92 { g.percent = (g.percent + 1).min(92); g.updated_ms = now_ms(); } } pub(crate) fn finish_done(&self, stage: impl Into, message: impl Into) { self.set(SelfUpdateState::Done, stage, 100, message); self.running.store(false, Ordering::SeqCst); } pub(crate) fn finish_failed(&self, stage: impl Into, message: impl Into) { self.set(SelfUpdateState::Failed, stage, 100, message); self.running.store(false, Ordering::SeqCst); } } #[derive(Debug, Clone, Serialize)] pub(crate) struct LibraryScanRow { pub(crate) kind: String, pub(crate) name: String, pub(crate) rel_path: String, pub(crate) ext: String, pub(crate) size: u64, #[serde(skip_serializing_if = "Option::is_none")] pub(crate) file_count: Option, pub(crate) modified_ms: Option, } #[derive(Debug, Clone, Serialize)] pub(crate) struct LibraryScanInfo { pub(crate) path: String, pub(crate) running: bool, pub(crate) scanned: usize, pub(crate) items: Vec, pub(crate) error: Option, pub(crate) started_ms: u64, pub(crate) updated_ms: u64, } #[derive(Debug, Default)] struct LibraryScanState { infos: HashMap, running: HashSet, } #[derive(Debug, Default)] pub(crate) struct LibraryScanStore { inner: Mutex, } impl LibraryScanStore { pub(crate) fn try_start(&self, path: String) -> bool { let mut g = self.inner.lock().unwrap_or_else(|e| e.into_inner()); if !g.running.insert(path.clone()) { return false; } let now = now_ms(); g.infos.insert( path.clone(), LibraryScanInfo { path, running: true, scanned: 0, items: Vec::new(), error: None, started_ms: now, updated_ms: now, }, ); true } pub(crate) fn push_item(&self, path: &str, item: LibraryScanRow, scanned: usize) { let mut g = self.inner.lock().unwrap_or_else(|e| e.into_inner()); if let Some(info) = g.infos.get_mut(path) { info.running = true; info.scanned = scanned; info.items.push(item); info.updated_ms = now_ms(); } } pub(crate) fn finish(&self, path: &str, mut items: Vec) { items.sort_by(|a, b| { b.modified_ms .cmp(&a.modified_ms) .then_with(|| a.name.cmp(&b.name)) }); let mut g = self.inner.lock().unwrap_or_else(|e| e.into_inner()); g.running.remove(path); let now = now_ms(); let started_ms = g.infos.get(path).map(|info| info.started_ms).unwrap_or(now); g.infos.insert( path.to_string(), LibraryScanInfo { path: path.to_string(), running: false, scanned: items.len(), items, error: None, started_ms, updated_ms: now, }, ); } pub(crate) fn finish_failed(&self, path: &str, error: String) { let mut g = self.inner.lock().unwrap_or_else(|e| e.into_inner()); g.running.remove(path); let now = now_ms(); let info = g.infos.entry(path.to_string()).or_insert(LibraryScanInfo { path: path.to_string(), running: false, scanned: 0, items: Vec::new(), error: None, started_ms: now, updated_ms: now, }); info.running = false; info.error = Some(error); info.updated_ms = now; } pub(crate) fn snapshot(&self, path: &str) -> Option { self.inner .lock() .unwrap_or_else(|e| e.into_inner()) .infos .get(path) .cloned() } } #[derive(Debug, Clone, Serialize)] pub(crate) struct UpdateScanRow { pub(crate) book_id: String, pub(crate) book_name: String, pub(crate) folder: String, pub(crate) local_total: usize, pub(crate) local_failed: usize, pub(crate) remote_total: usize, pub(crate) new_count: usize, pub(crate) has_update: bool, pub(crate) is_ignored: bool, } #[derive(Debug, Clone, Serialize)] pub(crate) struct UpdateScanInfo { pub(crate) running: bool, pub(crate) scanned: usize, pub(crate) total: usize, pub(crate) save_dir: String, pub(crate) updates: Vec, pub(crate) no_updates: Vec, pub(crate) error: Option, pub(crate) started_ms: u64, pub(crate) updated_ms: u64, } #[derive(Debug)] pub(crate) struct UpdateScanStore { running: AtomicBool, inner: Mutex, } impl Default for UpdateScanStore { fn default() -> Self { Self { running: AtomicBool::new(false), inner: Mutex::new(UpdateScanInfo { running: false, scanned: 0, total: 0, save_dir: String::new(), updates: Vec::new(), no_updates: Vec::new(), error: None, started_ms: 0, updated_ms: now_ms(), }), } } } impl UpdateScanStore { pub(crate) fn try_start(&self, save_dir: String) -> bool { if self .running .compare_exchange(false, true, Ordering::SeqCst, Ordering::SeqCst) .is_err() { return false; } let now = now_ms(); let mut g = self.inner.lock().unwrap_or_else(|e| e.into_inner()); *g = UpdateScanInfo { running: true, scanned: 0, total: 0, save_dir, updates: Vec::new(), no_updates: Vec::new(), error: None, started_ms: now, updated_ms: now, }; true } pub(crate) fn push_progress(&self, row: UpdateScanRow, scanned: usize, total: usize) { let mut g = self.inner.lock().unwrap_or_else(|e| e.into_inner()); g.running = true; g.scanned = scanned; g.total = total; g.updated_ms = now_ms(); if row.is_ignored || !row.has_update { g.no_updates.push(row); } else { g.updates.push(row); } } pub(crate) fn finish( &self, save_dir: String, updates: Vec, no_updates: Vec, ) { let total = updates.len() + no_updates.len(); let mut g = self.inner.lock().unwrap_or_else(|e| e.into_inner()); g.running = false; g.scanned = total; g.total = total; g.save_dir = save_dir; g.updates = updates; g.no_updates = no_updates; g.error = None; g.updated_ms = now_ms(); self.running.store(false, Ordering::SeqCst); } pub(crate) fn finish_failed(&self, error: String) { let mut g = self.inner.lock().unwrap_or_else(|e| e.into_inner()); g.running = false; g.error = Some(error); g.updated_ms = now_ms(); self.running.store(false, Ordering::SeqCst); } pub(crate) fn snapshot(&self) -> UpdateScanInfo { self.inner.lock().unwrap_or_else(|e| e.into_inner()).clone() } } #[derive(Clone)] pub(crate) struct AuthState { pub(crate) password_sha256: [u8; 32], pub(crate) session_secret: [u8; 32], cookie_secure: bool, login_attempts: Arc>>, } const SESSION_TTL_SECS: u64 = 24 * 60 * 60; const LOGIN_RATE_WINDOW_SECS: u64 = 1; const LOGIN_RATE_MAX_ATTEMPTS: usize = 5; const LOGIN_LOCK_AFTER_FAILURES: u32 = 10; const LOGIN_LOCK_SECS: u64 = 5 * 60; #[derive(Debug, Default)] struct LoginAttemptState { recent_attempts: VecDeque, failed_count: u32, locked_until: Option, } #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub(crate) enum LoginLimitDecision { Allowed, RateLimited { retry_after_secs: u64 }, Locked { retry_after_secs: u64 }, } impl AuthState { pub(crate) fn from_password(password: &str, cookie_secure: bool) -> Self { let mut h = Sha256::new(); h.update(password.as_bytes()); let out = h.finalize(); let mut password_sha256 = [0u8; 32]; password_sha256.copy_from_slice(&out); let nonce = Uuid::new_v4(); let now = now_secs(); let mut s = Sha256::new(); s.update(password_sha256); s.update(now.to_le_bytes()); s.update(nonce.as_bytes()); let secret = s.finalize(); let mut session_secret = [0u8; 32]; session_secret.copy_from_slice(&secret); Self { password_sha256, session_secret, cookie_secure, login_attempts: Arc::new(Mutex::new(HashMap::new())), } } pub(crate) fn check_login_allowed(&self, ip: IpAddr) -> LoginLimitDecision { let now = now_secs(); let mut attempts = self .login_attempts .lock() .unwrap_or_else(|e| e.into_inner()); let state = attempts.entry(ip).or_default(); if let Some(locked_until) = state.locked_until { if locked_until > now { return LoginLimitDecision::Locked { retry_after_secs: locked_until.saturating_sub(now).max(1), }; } state.locked_until = None; state.failed_count = 0; state.recent_attempts.clear(); } let cutoff = now.saturating_sub(LOGIN_RATE_WINDOW_SECS); while state .recent_attempts .front() .map(|ts| *ts <= cutoff) .unwrap_or(false) { state.recent_attempts.pop_front(); } if state.recent_attempts.len() >= LOGIN_RATE_MAX_ATTEMPTS { let oldest = *state.recent_attempts.front().unwrap_or(&now); let retry_after_secs = oldest .saturating_add(LOGIN_RATE_WINDOW_SECS) .saturating_sub(now) .max(1); return LoginLimitDecision::RateLimited { retry_after_secs }; } state.recent_attempts.push_back(now); LoginLimitDecision::Allowed } pub(crate) fn record_login_failure(&self, ip: IpAddr) -> Option { let now = now_secs(); let mut attempts = self .login_attempts .lock() .unwrap_or_else(|e| e.into_inner()); let state = attempts.entry(ip).or_default(); state.failed_count = state.failed_count.saturating_add(1); if state.failed_count >= LOGIN_LOCK_AFTER_FAILURES { let locked_until = now.saturating_add(LOGIN_LOCK_SECS); state.locked_until = Some(locked_until); state.recent_attempts.clear(); Some(LOGIN_LOCK_SECS) } else { None } } pub(crate) fn record_login_success(&self, ip: IpAddr) { let mut attempts = self .login_attempts .lock() .unwrap_or_else(|e| e.into_inner()); attempts.remove(&ip); } pub(crate) fn issue_session_token(&self) -> String { let exp = now_secs().saturating_add(SESSION_TTL_SECS); let nonce = Uuid::new_v4().simple().to_string(); let payload = format!("{exp}.{nonce}"); let sig = self.sign_payload(&payload); format!("{payload}.{sig}") } pub(crate) fn verify_session_token(&self, token: &str) -> bool { let mut parts = token.split('.'); let Some(exp_raw) = parts.next() else { return false; }; let Some(nonce_raw) = parts.next() else { return false; }; let Some(sig_raw) = parts.next() else { return false; }; if parts.next().is_some() { return false; } let Ok(exp) = exp_raw.parse::() else { return false; }; if now_secs() > exp { return false; } let payload = format!("{exp_raw}.{nonce_raw}"); let expected = self.sign_payload(&payload); constant_time_eq(sig_raw.as_bytes(), expected.as_bytes()) } pub(crate) fn session_ttl_secs(&self) -> u64 { SESSION_TTL_SECS } pub(crate) fn cookie_secure(&self) -> bool { self.cookie_secure } fn sign_payload(&self, payload: &str) -> String { let mut h = Sha256::new(); h.update(self.session_secret); h.update(payload.as_bytes()); hex::encode(h.finalize()) } } fn constant_time_eq(a: &[u8], b: &[u8]) -> bool { if a.len() != b.len() { return false; } let mut diff = 0u8; for (x, y) in a.iter().zip(b.iter()) { diff |= x ^ y; } diff == 0 } #[cfg(test)] mod auth_tests { use super::*; use std::net::IpAddr; #[test] fn login_rate_limit_allows_five_attempts_per_second() { let auth = AuthState::from_password("secret", false); let ip = IpAddr::from([127, 0, 0, 1]); for _ in 0..LOGIN_RATE_MAX_ATTEMPTS { assert_eq!(auth.check_login_allowed(ip), LoginLimitDecision::Allowed); } assert!(matches!( auth.check_login_allowed(ip), LoginLimitDecision::RateLimited { .. } )); } #[test] fn repeated_failures_lock_ip_and_success_resets_state() { let auth = AuthState::from_password("secret", false); let ip = IpAddr::from([127, 0, 0, 2]); for _ in 1..LOGIN_LOCK_AFTER_FAILURES { assert_eq!(auth.record_login_failure(ip), None); } assert_eq!(auth.record_login_failure(ip), Some(LOGIN_LOCK_SECS)); assert!(matches!( auth.check_login_allowed(ip), LoginLimitDecision::Locked { .. } )); auth.record_login_success(ip); assert_eq!(auth.check_login_allowed(ip), LoginLimitDecision::Allowed); } #[test] fn session_defaults_are_short_lived_and_secure_flag_is_configurable() { let insecure = AuthState::from_password("secret", false); let secure = AuthState::from_password("secret", true); assert_eq!(insecure.session_ttl_secs(), 24 * 60 * 60); assert!(!insecure.cookie_secure()); assert!(secure.cookie_secure()); assert!(secure.verify_session_token(&secure.issue_session_token())); } } pub(crate) const RECENT_DONE_JOB_RETENTION_MS: u64 = 2 * 60 * 60 * 1000; #[derive(Debug, Clone, Copy, Serialize)] #[serde(rename_all = "snake_case")] pub(crate) enum JobState { Queued, Running, Done, Failed, Canceled, } impl JobState { fn is_auto_prunable(self) -> bool { matches!(self, JobState::Done) } } #[derive(Debug, Clone, Serialize)] pub(crate) struct JobInfo { pub(crate) id: u64, pub(crate) book_id: String, pub(crate) title: Option, pub(crate) author: Option, pub(crate) state: JobState, pub(crate) progress: Option, pub(crate) message: Option, pub(crate) book_name_options: Option>, pub(crate) format_options: Option>, pub(crate) created_ms: u64, pub(crate) updated_ms: u64, } #[derive(Debug, Clone)] pub(crate) struct JobHandle { pub(crate) id: u64, pub(crate) cancel: Arc, } #[derive(Debug)] struct JobEntry { info: JobInfo, cancel: Arc, book_name_sender: Option>>, format_sender: Option>>, } #[derive(Debug, Default)] pub(crate) struct JobStore { next_id: AtomicU64, inner: Mutex>, } impl JobStore { pub(crate) fn create(&self, book_id: String) -> JobHandle { let id = self.next_id.fetch_add(1, Ordering::Relaxed) + 1; let now = now_ms(); let cancel = Arc::new(AtomicBool::new(false)); let info = JobInfo { id, book_id, title: None, author: None, state: JobState::Queued, progress: None, message: None, book_name_options: None, format_options: None, created_ms: now, updated_ms: now, }; let mut g = self.inner.lock().unwrap_or_else(|e| e.into_inner()); g.insert( id, JobEntry { info, cancel: cancel.clone(), book_name_sender: None, format_sender: None, }, ); JobHandle { id, cancel } } pub(crate) fn list(&self) -> Vec { let g = self.inner.lock().unwrap_or_else(|e| e.into_inner()); let mut v: Vec = g.values().map(|e| e.info.clone()).collect(); v.sort_by(|a, b| { b.updated_ms .cmp(&a.updated_ms) .then_with(|| b.id.cmp(&a.id)) }); v } pub(crate) fn prune_done_older_than(&self, retention_ms: u64) { let cutoff = now_ms().saturating_sub(retention_ms); let mut g = self.inner.lock().unwrap_or_else(|e| e.into_inner()); g.retain(|_, e| !e.info.state.is_auto_prunable() || e.info.updated_ms >= cutoff); } /// 返回当前处于 Queued 或 Running 状态的任务数量,用于并发限制。 pub(crate) fn count_active(&self) -> usize { let g = self.inner.lock().unwrap_or_else(|e| e.into_inner()); g.values() .filter(|e| matches!(e.info.state, JobState::Queued | JobState::Running)) .count() } pub(crate) fn set_running(&self, id: u64) { self.update(id, |j| { j.state = JobState::Running; j.message = None; j.book_name_options = None; j.format_options = None; }); } pub(crate) fn set_meta(&self, id: u64, title: Option, author: Option) { self.update(id, |j| { j.title = title; j.author = author; }); } pub(crate) fn set_progress(&self, id: u64, snap: ProgressSnapshot) { self.update(id, |j| { j.progress = Some(snap); }); } pub(crate) fn set_done(&self, id: u64) { self.update(id, |j| { j.state = JobState::Done; j.message = None; j.book_name_options = None; j.format_options = None; }); } pub(crate) fn set_failed(&self, id: u64, msg: String) { self.update(id, |j| { j.state = JobState::Failed; j.message = Some(msg); j.book_name_options = None; j.format_options = None; }); } pub(crate) fn request_cancel(&self, id: u64) -> bool { let mut g = self.inner.lock().unwrap_or_else(|e| e.into_inner()); let Some(e) = g.get_mut(&id) else { return false; }; e.cancel.store(true, Ordering::Relaxed); e.info.state = JobState::Canceled; e.info.message = Some("cancel requested".to_string()); if let Some(tx) = e.book_name_sender.take() { let _ = tx.send(None); } if let Some(tx) = e.format_sender.take() { let _ = tx.send(None); } e.info.book_name_options = None; e.info.format_options = None; e.info.updated_ms = now_ms(); true } pub(crate) fn request_cancel_and_remove(&self, id: u64) -> bool { let mut g = self.inner.lock().unwrap_or_else(|e| e.into_inner()); let Some(mut e) = g.remove(&id) else { return false; }; e.cancel.store(true, Ordering::Relaxed); if let Some(tx) = e.book_name_sender.take() { let _ = tx.send(None); } if let Some(tx) = e.format_sender.take() { let _ = tx.send(None); } true } pub(crate) fn remove(&self, id: u64) -> bool { let mut g = self.inner.lock().unwrap_or_else(|e| e.into_inner()); let Some(mut e) = g.remove(&id) else { return false; }; if let Some(tx) = e.book_name_sender.take() { let _ = tx.send(None); } if let Some(tx) = e.format_sender.take() { let _ = tx.send(None); } true } pub(crate) fn set_book_name_options( &self, id: u64, options: Vec, sender: std::sync::mpsc::Sender>, ) { let mut g = self.inner.lock().unwrap_or_else(|e| e.into_inner()); let Some(e) = g.get_mut(&id) else { return; }; e.info.book_name_options = Some(options); e.info.message = Some("等待选择书名".to_string()); e.book_name_sender = Some(sender); e.info.updated_ms = now_ms(); } pub(crate) fn submit_book_name_choice(&self, id: u64, choice: Option) -> bool { let mut g = self.inner.lock().unwrap_or_else(|e| e.into_inner()); let Some(e) = g.get_mut(&id) else { return false; }; if let Some(tx) = e.book_name_sender.take() { let _ = tx.send(choice); e.info.book_name_options = None; e.info.message = None; e.info.updated_ms = now_ms(); return true; } false } pub(crate) fn set_format_options( &self, id: u64, options: Vec, sender: std::sync::mpsc::Sender>, ) { let mut g = self.inner.lock().unwrap_or_else(|e| e.into_inner()); let Some(e) = g.get_mut(&id) else { return; }; e.info.format_options = Some(options); e.info.message = Some("等待选择输出格式".to_string()); e.format_sender = Some(sender); e.info.updated_ms = now_ms(); } pub(crate) fn submit_format_choice(&self, id: u64, choice: Option) -> bool { let mut g = self.inner.lock().unwrap_or_else(|e| e.into_inner()); let Some(e) = g.get_mut(&id) else { return false; }; if let Some(tx) = e.format_sender.take() { let _ = tx.send(choice); e.info.format_options = None; e.info.message = None; e.info.updated_ms = now_ms(); return true; } false } fn update(&self, id: u64, f: F) { let mut g = self.inner.lock().unwrap_or_else(|e| e.into_inner()); let Some(e) = g.get_mut(&id) else { return; }; f(&mut e.info); e.info.updated_ms = now_ms(); } } fn now_ms() -> u64 { SystemTime::now() .duration_since(UNIX_EPOCH) .unwrap_or_default() .as_millis() as u64 } fn now_secs() -> u64 { SystemTime::now() .duration_since(UNIX_EPOCH) .unwrap_or_default() .as_secs() }